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::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
53/// The result type returned by a [`FrontendHeartbeatExtension`].
54pub type FrontendHeartbeatExtensionResult<T> = std::result::Result<T, BoxedError>;
55
56/// An extension to frontend heartbeat requests, responses, and lifecycle events.
57///
58/// [`FrontendHeartbeatExtension::request_extensions`] is called for every heartbeat. A failed call
59/// is isolated from the base heartbeat and from other extensions.
60/// [`FrontendHeartbeatExtension::connected`] is called once for every successfully established
61/// connection generation, including reconnects; implementations must therefore be idempotent.
62/// [`FrontendHeartbeatExtension::shutdown`] is called once during task shutdown after heartbeat
63/// I/O has stopped.
64#[async_trait]
65pub trait FrontendHeartbeatExtension: Send + Sync {
66    /// Returns the stable name used to make registration idempotent.
67    fn name(&self) -> &str;
68
69    /// Generates request extensions for one heartbeat.
70    async fn request_extensions(
71        &self,
72    ) -> FrontendHeartbeatExtensionResult<HashMap<String, Vec<u8>>> {
73        Ok(HashMap::new())
74    }
75
76    /// Returns a handler to insert into the heartbeat response handler chain.
77    ///
78    /// Errors and [`common_meta::heartbeat::handler::HandleControl::Done`] are isolated to this
79    /// extension so they cannot skip later extensions or mandatory OSS handlers.
80    fn response_handler(&self) -> Option<HeartbeatResponseHandlerRef> {
81        None
82    }
83
84    /// Notifies the extension that a heartbeat connection generation is ready.
85    async fn connected(&self, _generation: u64) -> FrontendHeartbeatExtensionResult<()> {
86        Ok(())
87    }
88
89    /// Stops and joins background work owned by the extension.
90    async fn shutdown(&self) -> FrontendHeartbeatExtensionResult<()> {
91        Ok(())
92    }
93}
94
95/// A shareable, ordered registry of frontend heartbeat extensions.
96#[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    /// Registers an extension without replacing an existing extension of the same name.
109    ///
110    /// Returns `true` for a new registration and `false` for an idempotent duplicate.
111    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    /// Returns the registered extensions in registration order.
122    pub fn extensions(&self) -> Vec<Arc<dyn FrontendHeartbeatExtension>> {
123        self.inner.lock().unwrap().extensions.clone()
124    }
125
126    /// Returns response handlers in extension registration order.
127    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    /// Returns the number of registered extensions.
135    pub fn len(&self) -> usize {
136        self.inner.lock().unwrap().extensions.len()
137    }
138
139    /// Returns whether no extension is registered.
140    pub fn is_empty(&self) -> bool {
141        self.len() == 0
142    }
143}
144
145/// Builds the frontend heartbeat response handler chain.
146///
147/// Mailbox parsing always runs first, followed by extension handlers, suspension state handling,
148/// and cache invalidation.
149pub 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                // Metasrv calculates the frontend id by hashing this reachable address.
407                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            // TODO(zyy17): Remove these deprecated fields when the deprecated fields are removed from the proto.
537            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/// The frontend task that sends [`HeartbeatRequest`] values to metasrv in the background.
561#[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    /// Installs the extensions registered before heartbeat startup.
614    pub fn with_extensions(mut self, extensions: FrontendHeartbeatExtensions) -> Self {
615        self.runner.extensions = extensions;
616        self
617    }
618
619    /// Establishes the initial heartbeat connection and starts its background supervisor.
620    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    /// Cancels and joins heartbeat I/O and all registered extension lifecycles.
663    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        // Wait for a concurrently running handshake or connected callback to observe cancellation.
671        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 internal grpc is configured, use its address as the peer address
700    // otherwise use the public grpc address, because peer address only promises to be reachable
701    // by other components, it doesn't matter whether it's internal or external
702    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}