Skip to main content

frontend/
heartbeat.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
15#[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
54/// The result type returned by a [`FrontendHeartbeatExtension`].
55pub type FrontendHeartbeatExtensionResult<T> = std::result::Result<T, BoxedError>;
56
57/// An extension to frontend heartbeat requests, responses, and lifecycle events.
58///
59/// [`FrontendHeartbeatExtension::request_extensions`] is called for every heartbeat. A failed call
60/// is isolated from the base heartbeat and from other extensions.
61/// [`FrontendHeartbeatExtension::connected`] is called once for every successfully established
62/// connection generation, including reconnects; implementations must therefore be idempotent.
63/// [`FrontendHeartbeatExtension::shutdown`] is called once during task shutdown after heartbeat
64/// I/O has stopped.
65#[async_trait]
66pub trait FrontendHeartbeatExtension: Send + Sync {
67    /// Returns the stable name used to make registration idempotent.
68    fn name(&self) -> &str;
69
70    /// Generates request extensions for one heartbeat.
71    async fn request_extensions(
72        &self,
73    ) -> FrontendHeartbeatExtensionResult<HashMap<String, Vec<u8>>> {
74        Ok(HashMap::new())
75    }
76
77    /// Returns a handler to insert into the heartbeat response handler chain.
78    ///
79    /// Errors and [`common_meta::heartbeat::handler::HandleControl::Done`] are isolated to this
80    /// extension so they cannot skip later extensions or mandatory OSS handlers.
81    fn response_handler(&self) -> Option<HeartbeatResponseHandlerRef> {
82        None
83    }
84
85    /// Notifies the extension that a heartbeat connection generation is ready.
86    async fn connected(&self, _generation: u64) -> FrontendHeartbeatExtensionResult<()> {
87        Ok(())
88    }
89
90    /// Stops and joins background work owned by the extension.
91    async fn shutdown(&self) -> FrontendHeartbeatExtensionResult<()> {
92        Ok(())
93    }
94}
95
96/// A shareable, ordered registry of frontend heartbeat extensions.
97///
98/// Registration is available until [`FrontendHeartbeatExtensions::freeze`] is called. After that,
99/// membership is immutable so every heartbeat consumer observes the same extensions.
100#[derive(Clone, Default)]
101pub struct FrontendHeartbeatExtensions {
102    inner: Arc<StdMutex<FrontendHeartbeatExtensionsInner>>,
103}
104
105/// The reason a frontend heartbeat extension was not registered.
106#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107pub enum RegistrationError {
108    /// An extension with the same name is already registered.
109    Duplicate,
110    /// The registry is frozen and no longer accepts registrations.
111    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    /// Registers an extension without replacing an existing extension of the same name.
138    ///
139    /// Returns `true` for a new registration and `false` for an idempotent duplicate or when the
140    /// registry is frozen. Use [`FrontendHeartbeatExtensions::try_register`] when the rejection
141    /// reason must not be discarded.
142    pub fn register(&self, extension: Arc<dyn FrontendHeartbeatExtension>) -> bool {
143        self.try_register(extension).is_ok()
144    }
145
146    /// Registers an extension and returns the reason when registration is rejected.
147    ///
148    /// If the registry is already frozen, this method returns [`RegistrationError::Frozen`]
149    /// without calling [`FrontendHeartbeatExtension::name`]. If freezing races with name
150    /// evaluation, the registry state is checked again before insertion.
151    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    /// Prevents further registrations while preserving the current registration order.
172    pub fn freeze(&self) {
173        self.lock().frozen = true;
174    }
175
176    /// Returns the registered extensions in registration order.
177    pub fn extensions(&self) -> Vec<Arc<dyn FrontendHeartbeatExtension>> {
178        self.lock().extensions.clone()
179    }
180
181    /// Returns response handlers in extension registration order.
182    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    /// Returns the number of registered extensions.
190    pub fn len(&self) -> usize {
191        self.lock().extensions.len()
192    }
193
194    /// Returns whether no extension is registered.
195    pub fn is_empty(&self) -> bool {
196        self.len() == 0
197    }
198}
199
200/// Builds the frontend heartbeat response handler chain.
201///
202/// Mailbox parsing always runs first, followed by extension handlers, suspension state handling,
203/// and cache invalidation.
204pub 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                // Metasrv calculates the frontend id by hashing this reachable address.
462                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            // TODO(zyy17): Remove these deprecated fields when the deprecated fields are removed from the proto.
592            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/// The frontend task that sends [`HeartbeatRequest`] values to metasrv in the background.
616#[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    /// Installs the extensions registered before heartbeat startup.
669    pub fn with_extensions(mut self, extensions: FrontendHeartbeatExtensions) -> Self {
670        self.runner.extensions = extensions;
671        self
672    }
673
674    /// Establishes the initial heartbeat connection and starts its background supervisor.
675    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    /// Cancels and joins heartbeat I/O and all registered extension lifecycles.
718    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        // Wait for a concurrently running handshake or connected callback to observe cancellation.
726        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 internal grpc is configured, use its address as the peer address
755    // otherwise use the public grpc address, because peer address only promises to be reachable
756    // by other components, it doesn't matter whether it's internal or external
757    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}