feat: make frontend heartbeat extensible and lifecycle-safe (#8726)

* feat: make frontend heartbeat extensible and lifecycle-safe

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix: isolate heartbeat extension response handlers

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix: cancel in-flight heartbeat response handling

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test: cover heartbeat wire compatibility

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix: clean up failed heartbeat startup

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix: address frontend heartbeat review feedback

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

---------

Signed-off-by: jeremyhi <fengjiachun@gmail.com>
(cherry picked from commit 6ffa27bc8e78b5e11b2db3d2a898438446b298be)
This commit is contained in:
jeremyhi
2026-08-05 12:38:54 +08:00
committed by discord9
parent 2f5e97850e
commit 424e53d841
9 changed files with 1817 additions and 194 deletions
Generated
+12 -12
View File
@@ -1623,7 +1623,7 @@ dependencies = [
"maybe-owned",
"rustix 1.0.7",
"rustix-linux-procfs",
"windows-sys 0.61.2",
"windows-sys 0.60.2",
"winx",
]
@@ -2259,7 +2259,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "117725a109d387c937a1533ce01b450cbde6b88abceea8473c4d7a85853cda3c"
dependencies = [
"lazy_static",
"windows-sys 0.59.0",
"windows-sys 0.48.0",
]
[[package]]
@@ -6041,7 +6041,7 @@ dependencies = [
[[package]]
name = "greptime-proto"
version = "0.1.0"
source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=5adb8a1637abe87bcbb455d32136eee5d538cc49#5adb8a1637abe87bcbb455d32136eee5d538cc49"
source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=032510ded061277b7d1bbb29d4046abb5b1cbf4b#032510ded061277b7d1bbb29d4046abb5b1cbf4b"
dependencies = [
"prost 0.14.1",
"prost-types 0.14.1",
@@ -6587,7 +6587,7 @@ dependencies = [
"libc",
"percent-encoding",
"pin-project-lite",
"socket2 0.6.4",
"socket2 0.5.10",
"tokio",
"tower-service",
"tracing",
@@ -9095,7 +9095,7 @@ version = "0.50.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5"
dependencies = [
"windows-sys 0.61.2",
"windows-sys 0.60.2",
]
[[package]]
@@ -11284,7 +11284,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ac6c3320f9abac597dcbc668774ef006702672474aad53c6d596b62e487b40b1"
dependencies = [
"heck 0.5.0",
"itertools 0.14.0",
"itertools 0.10.5",
"log",
"multimap",
"once_cell",
@@ -11332,7 +11332,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d"
dependencies = [
"anyhow",
"itertools 0.14.0",
"itertools 0.10.5",
"proc-macro2",
"quote",
"syn 2.0.117",
@@ -11345,7 +11345,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9120690fafc389a67ba3803df527d0ec9cbbc9cc45e4cc20b332996dfb672425"
dependencies = [
"anyhow",
"itertools 0.14.0",
"itertools 0.10.5",
"proc-macro2",
"quote",
"syn 2.0.117",
@@ -12800,7 +12800,7 @@ dependencies = [
"security-framework 3.7.0",
"security-framework-sys",
"webpki-root-certs",
"windows-sys 0.61.2",
"windows-sys 0.60.2",
]
[[package]]
@@ -13712,7 +13712,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "52d1cfed4120b4d927bf7c0f86d2087a4a7d6027c906d9f9d525a80573b9be51"
dependencies = [
"libc",
"windows-sys 0.61.2",
"windows-sys 0.60.2",
]
[[package]]
@@ -14706,7 +14706,7 @@ dependencies = [
"getrandom 0.3.4",
"once_cell",
"rustix 1.0.7",
"windows-sys 0.61.2",
"windows-sys 0.60.2",
]
[[package]]
@@ -16459,7 +16459,7 @@ version = "0.1.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cf221c93e13a30d793f7645a0e7762c55d169dbb0a49671918a2319d289b10bb"
dependencies = [
"windows-sys 0.59.0",
"windows-sys 0.48.0",
]
[[package]]
+1 -1
View File
@@ -159,7 +159,7 @@ fs2 = "0.4"
fst = "0.4.7"
futures = "0.3"
futures-util = "0.3"
greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "5adb8a1637abe87bcbb455d32136eee5d538cc49" }
greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "032510ded061277b7d1bbb29d4046abb5b1cbf4b" }
hex = "0.4"
http = "1"
humantime = "2.1"
+31 -14
View File
@@ -32,10 +32,6 @@ use common_base::Plugins;
use common_config::{Configurable, DEFAULT_DATA_HOME};
use common_error::ext::BoxedError;
use common_meta::cache::{CacheRegistryBuilder, LayeredCacheRegistryBuilder};
use common_meta::heartbeat::handler::HandlerGroupExecutor;
use common_meta::heartbeat::handler::invalidate_table_cache::InvalidateCacheHandler;
use common_meta::heartbeat::handler::parse_mailbox_message::ParseMailboxMessageHandler;
use common_meta::heartbeat::handler::suspend::SuspendHandler;
use common_query::prelude::set_default_prefix;
use common_stat::ResourceStatImpl;
use common_telemetry::info;
@@ -43,7 +39,9 @@ use common_telemetry::logging::{DEFAULT_LOGGING_DIR, TracingOptions};
use common_time::timezone::set_default_timezone;
use common_version::{short_version, verbose_version};
use frontend::frontend::Frontend;
use frontend::heartbeat::HeartbeatTask;
use frontend::heartbeat::{
FrontendHeartbeatExtensions, HeartbeatTask, heartbeat_response_handler_executor,
};
use frontend::instance::builder::FrontendBuilder;
use frontend::server::Services;
use meta_client::{MetaClientOptions, MetaClientRef, MetaClientType};
@@ -497,10 +495,21 @@ impl StartCommand {
.await
.context(error::StartFrontendSnafu)?;
let heartbeat_task = Some(create_heartbeat_task(&opts, meta_client, &instance));
let instance = Arc::new(instance);
plugins::setup_frontend_heartbeat_extensions(&mut plugins, &plugin_opts, &instance)
.await
.context(error::StartFrontendSnafu)?;
let heartbeat_extensions = plugins
.get::<FrontendHeartbeatExtensions>()
.unwrap_or_default();
let heartbeat_task = Some(create_heartbeat_task_with_extensions(
&opts,
meta_client,
&instance,
heartbeat_extensions,
));
let servers = Services::new(opts, instance.clone(), plugins)
.build()
.context(error::StartFrontendSnafu)?;
@@ -520,13 +529,20 @@ pub fn create_heartbeat_task(
meta_client: MetaClientRef,
instance: &frontend::instance::Instance,
) -> HeartbeatTask {
let executor = Arc::new(HandlerGroupExecutor::new(vec![
Arc::new(ParseMailboxMessageHandler),
Arc::new(SuspendHandler::new(instance.suspend_state())),
Arc::new(InvalidateCacheHandler::new(
instance.cache_invalidator().clone(),
)),
]));
create_heartbeat_task_with_extensions(options, meta_client, instance, Default::default())
}
fn create_heartbeat_task_with_extensions(
options: &frontend::frontend::FrontendOptions,
meta_client: MetaClientRef,
instance: &frontend::instance::Instance,
extensions: FrontendHeartbeatExtensions,
) -> HeartbeatTask {
let executor = heartbeat_response_handler_executor(
&extensions,
instance.suspend_state(),
instance.cache_invalidator().clone(),
);
let stat = {
let mut stat = ResourceStatImpl::default();
@@ -541,6 +557,7 @@ pub fn create_heartbeat_task(
executor,
stat,
)
.with_extensions(extensions)
}
#[cfg(test)]
+154 -7
View File
@@ -130,17 +130,26 @@ pub struct Frontend {
impl Frontend {
pub async fn start(&mut self) -> Result<()> {
if let Some(t) = &self.heartbeat_task {
t.start().await?;
if let Some(t) = &self.heartbeat_task
&& let Err(error) = t.start().await
{
t.shutdown().await;
return Err(error);
}
self.servers
.start_all()
.await
.context(error::StartServerSnafu)
if let Err(source) = self.servers.start_all().await {
if let Some(t) = &self.heartbeat_task {
t.shutdown().await;
}
return Err(source).context(error::StartServerSnafu);
}
Ok(())
}
pub async fn shutdown(&mut self) -> Result<()> {
if let Some(t) = &self.heartbeat_task {
t.shutdown().await;
}
self.servers
.shutdown_all()
.await
@@ -154,7 +163,9 @@ impl Frontend {
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicBool, Ordering};
use std::any::Any;
use std::net::SocketAddr;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::time::Duration;
use api::v1::meta::heartbeat_server::HeartbeatServer;
@@ -181,6 +192,7 @@ mod tests {
use servers::grpc::{FlightCompression, GRPC_SERVER};
use servers::http::HTTP_SERVER;
use servers::http::result::greptime_result_v1::GreptimedbV1Response;
use servers::server::Server;
use tokio::sync::mpsc;
use tonic::codec::CompressionEncoding;
use tonic::codegen::tokio_stream::StreamExt;
@@ -188,6 +200,9 @@ mod tests {
use tonic::{Request, Response, Status, Streaming};
use super::*;
use crate::heartbeat::{
FrontendHeartbeatExtension, FrontendHeartbeatExtensionResult, FrontendHeartbeatExtensions,
};
use crate::instance::builder::FrontendBuilder;
use crate::server::Services;
@@ -209,6 +224,46 @@ mod tests {
struct SuspendableHeartbeatServer {
suspend: Arc<AtomicBool>,
fail_heartbeat: bool,
}
struct FailingServer;
struct ShutdownTrackingExtension {
shutdown_calls: AtomicUsize,
}
#[async_trait]
impl FrontendHeartbeatExtension for ShutdownTrackingExtension {
fn name(&self) -> &str {
"shutdown-tracking"
}
async fn shutdown(&self) -> FrontendHeartbeatExtensionResult<()> {
self.shutdown_calls.fetch_add(1, Ordering::AcqRel);
Ok(())
}
}
#[async_trait]
impl Server for FailingServer {
async fn shutdown(&self) -> servers::error::Result<()> {
Ok(())
}
async fn start(&mut self, _listening: SocketAddr) -> servers::error::Result<()> {
Err(servers::error::Error::Internal {
err_msg: "mock server start failure".to_string(),
})
}
fn name(&self) -> &str {
"FAILING_SERVER"
}
fn as_any(&self) -> &dyn Any {
self
}
}
#[async_trait]
@@ -219,6 +274,10 @@ mod tests {
&self,
request: Request<Streaming<HeartbeatRequest>>,
) -> std::result::Result<Response<Self::HeartbeatStream>, Status> {
if self.fail_heartbeat {
return Err(Status::unavailable("mock initial heartbeat failure"));
}
let (tx, rx) = mpsc::channel(4);
common_runtime::spawn_global({
@@ -358,6 +417,93 @@ mod tests {
Ok(frontend)
}
#[tokio::test]
async fn test_server_start_failure_shuts_down_heartbeat() {
let meta_client_options = MetaClientOptions {
metasrv_addrs: vec!["localhost:0".to_string()],
..Default::default()
};
let options = FrontendOptions {
meta_client: Some(meta_client_options.clone()),
..Default::default()
};
let heartbeat_server = Arc::new(SuspendableHeartbeatServer {
suspend: Arc::new(AtomicBool::new(false)),
fail_heartbeat: false,
});
let meta_client = create_meta_client(&meta_client_options, heartbeat_server).await;
let instance = Arc::new(
FrontendBuilder::new_test(&options, meta_client.clone())
.try_build()
.await
.unwrap(),
);
let heartbeat_task = HeartbeatTask::new(
instance.frontend_peer_addr().to_string(),
&options,
meta_client,
Arc::new(HandlerGroupExecutor::new(vec![])),
Arc::new(ResourceStatImpl::default()),
);
let heartbeat_probe = heartbeat_task.clone();
let servers = ServerHandlers::default();
servers.insert((Box::new(FailingServer), "127.0.0.1:0".parse().unwrap()));
let mut frontend = Frontend {
instance,
servers,
heartbeat_task: Some(heartbeat_task),
};
assert!(frontend.start().await.is_err());
assert!(heartbeat_probe.is_shutdown());
}
#[tokio::test]
async fn test_heartbeat_start_failure_shuts_down_extensions() {
let meta_client_options = MetaClientOptions {
metasrv_addrs: vec!["localhost:0".to_string()],
..Default::default()
};
let options = FrontendOptions {
meta_client: Some(meta_client_options.clone()),
..Default::default()
};
let heartbeat_server = Arc::new(SuspendableHeartbeatServer {
suspend: Arc::new(AtomicBool::new(false)),
fail_heartbeat: true,
});
let meta_client = create_meta_client(&meta_client_options, heartbeat_server).await;
let instance = Arc::new(
FrontendBuilder::new_test(&options, meta_client.clone())
.try_build()
.await
.unwrap(),
);
let extension = Arc::new(ShutdownTrackingExtension {
shutdown_calls: AtomicUsize::new(0),
});
let extensions = FrontendHeartbeatExtensions::default();
assert!(extensions.register(extension.clone()));
let heartbeat_task = HeartbeatTask::new(
instance.frontend_peer_addr().to_string(),
&options,
meta_client,
Arc::new(HandlerGroupExecutor::new(vec![])),
Arc::new(ResourceStatImpl::default()),
)
.with_extensions(extensions);
let heartbeat_probe = heartbeat_task.clone();
let mut frontend = Frontend {
instance,
servers: ServerHandlers::default(),
heartbeat_task: Some(heartbeat_task),
};
assert!(frontend.start().await.is_err());
assert!(heartbeat_probe.is_shutdown());
assert_eq!(extension.shutdown_calls.load(Ordering::Acquire), 1);
}
async fn verify_suspend_state_by_http(
frontend: &Frontend,
expected: std::result::Result<&str, (StatusCode, &str)>,
@@ -454,6 +600,7 @@ mod tests {
let server = Arc::new(SuspendableHeartbeatServer {
suspend: Arc::new(AtomicBool::new(false)),
fail_heartbeat: false,
});
let meta_client = create_meta_client(&meta_client_options, server.clone()).await;
let frontend = create_frontend(&options, meta_client).await?;
+578 -153
View File
@@ -15,119 +15,481 @@
#[cfg(test)]
mod tests;
use std::sync::Arc;
use std::collections::{HashMap, HashSet};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex as StdMutex};
use api::v1::meta::heartbeat_request::NodeWorkloads;
use api::v1::meta::{FrontendWorkloads, HeartbeatRequest, NodeInfo, Peer};
use api::v1::meta::{FrontendWorkloads, HeartbeatRequest, HeartbeatResponse, NodeInfo, Peer};
use async_trait::async_trait;
use common_error::ext::BoxedError;
use common_meta::cache_invalidator::CacheInvalidatorRef;
use common_meta::datanode::EnvVars;
use common_meta::heartbeat::handler::invalidate_table_cache::InvalidateCacheHandler;
use common_meta::heartbeat::handler::parse_mailbox_message::ParseMailboxMessageHandler;
use common_meta::heartbeat::handler::suspend::SuspendHandler;
use common_meta::heartbeat::handler::{
HeartbeatResponseHandlerContext, HeartbeatResponseHandlerExecutorRef,
HandlerGroupExecutor, HeartbeatResponseHandlerContext, HeartbeatResponseHandlerExecutorRef,
HeartbeatResponseHandlerRef,
};
use common_meta::heartbeat::mailbox::{HeartbeatMailbox, MailboxRef, OutgoingMessage};
use common_meta::heartbeat::utils::outgoing_message_to_mailbox_message;
use common_stat::ResourceStatRef;
use common_telemetry::{debug, error, info, warn};
use meta_client::client::heartbeat::HeartbeatConfig;
use meta_client::client::{HeartbeatSender, HeartbeatStream, MetaClient};
use servers::addrs;
use snafu::ResultExt;
use tokio::sync::mpsc;
use tokio::sync::mpsc::Receiver;
use tokio::sync::{Mutex, mpsc};
use tokio::time::{Duration, Instant};
use tokio_util::sync::CancellationToken;
use crate::error;
use crate::error::Result;
use crate::frontend::FrontendOptions;
use crate::metrics::{HEARTBEAT_RECV_COUNT, HEARTBEAT_SENT_COUNT};
/// The frontend heartbeat task which sending `[HeartbeatRequest]` to Metasrv periodically in background.
/// The result type returned by a [`FrontendHeartbeatExtension`].
pub type FrontendHeartbeatExtensionResult<T> = std::result::Result<T, BoxedError>;
/// An extension to frontend heartbeat requests, responses, and lifecycle events.
///
/// [`FrontendHeartbeatExtension::request_extensions`] is called for every heartbeat. A failed call
/// is isolated from the base heartbeat and from other extensions.
/// [`FrontendHeartbeatExtension::connected`] is called once for every successfully established
/// connection generation, including reconnects; implementations must therefore be idempotent.
/// [`FrontendHeartbeatExtension::shutdown`] is called once during task shutdown after heartbeat
/// I/O has stopped.
#[async_trait]
pub trait FrontendHeartbeatExtension: Send + Sync {
/// Returns the stable name used to make registration idempotent.
fn name(&self) -> &str;
/// Generates request extensions for one heartbeat.
async fn request_extensions(
&self,
) -> FrontendHeartbeatExtensionResult<HashMap<String, Vec<u8>>> {
Ok(HashMap::new())
}
/// Returns a handler to insert into the heartbeat response handler chain.
///
/// Errors and [`common_meta::heartbeat::handler::HandleControl::Done`] are isolated to this
/// extension so they cannot skip later extensions or mandatory OSS handlers.
fn response_handler(&self) -> Option<HeartbeatResponseHandlerRef> {
None
}
/// Notifies the extension that a heartbeat connection generation is ready.
async fn connected(&self, _generation: u64) -> FrontendHeartbeatExtensionResult<()> {
Ok(())
}
/// Stops and joins background work owned by the extension.
async fn shutdown(&self) -> FrontendHeartbeatExtensionResult<()> {
Ok(())
}
}
/// A shareable, ordered registry of frontend heartbeat extensions.
#[derive(Clone, Default)]
pub struct FrontendHeartbeatExtensions {
inner: Arc<StdMutex<FrontendHeartbeatExtensionsInner>>,
}
#[derive(Default)]
struct FrontendHeartbeatExtensionsInner {
names: HashSet<String>,
extensions: Vec<Arc<dyn FrontendHeartbeatExtension>>,
}
impl FrontendHeartbeatExtensions {
/// Registers an extension without replacing an existing extension of the same name.
///
/// Returns `true` for a new registration and `false` for an idempotent duplicate.
pub fn register(&self, extension: Arc<dyn FrontendHeartbeatExtension>) -> bool {
let mut inner = self.inner.lock().unwrap();
let name = extension.name().to_string();
if !inner.names.insert(name) {
return false;
}
inner.extensions.push(extension);
true
}
/// Returns the registered extensions in registration order.
pub fn extensions(&self) -> Vec<Arc<dyn FrontendHeartbeatExtension>> {
self.inner.lock().unwrap().extensions.clone()
}
/// Returns response handlers in extension registration order.
pub fn response_handlers(&self) -> Vec<HeartbeatResponseHandlerRef> {
self.extensions()
.into_iter()
.filter_map(|extension| extension.response_handler())
.collect()
}
/// Returns the number of registered extensions.
pub fn len(&self) -> usize {
self.inner.lock().unwrap().extensions.len()
}
/// Returns whether no extension is registered.
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
/// Builds the frontend heartbeat response handler chain.
///
/// Mailbox parsing always runs first, followed by extension handlers, suspension state handling,
/// and cache invalidation.
pub fn heartbeat_response_handler_executor(
extensions: &FrontendHeartbeatExtensions,
suspend_state: Arc<AtomicBool>,
cache_invalidator: CacheInvalidatorRef,
) -> HeartbeatResponseHandlerExecutorRef {
let mut handlers: Vec<HeartbeatResponseHandlerRef> = vec![Arc::new(ParseMailboxMessageHandler)];
handlers.extend(extensions.response_handlers().into_iter().map(|handler| {
Arc::new(IsolatedHeartbeatResponseHandler(handler)) as HeartbeatResponseHandlerRef
}));
handlers.extend([
Arc::new(SuspendHandler::new(suspend_state)) as HeartbeatResponseHandlerRef,
Arc::new(InvalidateCacheHandler::new(cache_invalidator)),
]);
Arc::new(HandlerGroupExecutor::new(handlers))
}
struct IsolatedHeartbeatResponseHandler(HeartbeatResponseHandlerRef);
#[async_trait]
impl common_meta::heartbeat::handler::HeartbeatResponseHandler
for IsolatedHeartbeatResponseHandler
{
fn is_acceptable(&self, _ctx: &HeartbeatResponseHandlerContext) -> bool {
true
}
async fn handle(
&self,
ctx: &mut HeartbeatResponseHandlerContext,
) -> common_meta::error::Result<common_meta::heartbeat::handler::HandleControl> {
use common_meta::heartbeat::handler::HandleControl;
if self.0.is_acceptable(ctx)
&& let Err(error) = self.0.handle(ctx).await
{
error!(error; "Heartbeat extension response handler failed");
}
Ok(HandleControl::Continue)
}
}
#[async_trait]
trait HeartbeatConnector: Send + Sync {
async fn connect(&self) -> Result<HeartbeatConnection>;
}
#[async_trait]
trait HeartbeatRequestSender: Send + Sync {
async fn send(&self, request: HeartbeatRequest) -> FrontendHeartbeatExtensionResult<()>;
}
#[async_trait]
trait HeartbeatResponseStream: Send {
async fn message(&mut self) -> FrontendHeartbeatExtensionResult<Option<HeartbeatResponse>>;
}
struct HeartbeatConnection {
sender: Arc<dyn HeartbeatRequestSender>,
stream: Box<dyn HeartbeatResponseStream>,
config: HeartbeatConfig,
}
struct MetaHeartbeatConnector {
client: Arc<MetaClient>,
}
#[async_trait]
impl HeartbeatConnector for MetaHeartbeatConnector {
async fn connect(&self) -> Result<HeartbeatConnection> {
let (sender, stream, config) = self
.client
.heartbeat()
.await
.context(error::CreateMetaHeartbeatStreamSnafu)?;
Ok(HeartbeatConnection {
sender: Arc::new(MetaHeartbeatSender(sender)),
stream: Box::new(MetaHeartbeatStream(stream)),
config,
})
}
}
struct MetaHeartbeatSender(HeartbeatSender);
#[async_trait]
impl HeartbeatRequestSender for MetaHeartbeatSender {
async fn send(&self, request: HeartbeatRequest) -> FrontendHeartbeatExtensionResult<()> {
self.0.send(request).await.map_err(BoxedError::new)
}
}
struct MetaHeartbeatStream(HeartbeatStream);
#[async_trait]
impl HeartbeatResponseStream for MetaHeartbeatStream {
async fn message(&mut self) -> FrontendHeartbeatExtensionResult<Option<HeartbeatResponse>> {
self.0.message().await.map_err(BoxedError::new)
}
}
#[derive(Clone)]
pub struct HeartbeatTask {
struct HeartbeatRunner {
peer_addr: String,
meta_client: Arc<MetaClient>,
connector: Arc<dyn HeartbeatConnector>,
resp_handler_executor: HeartbeatResponseHandlerExecutorRef,
start_time_ms: u64,
resource_stat: ResourceStatRef,
env_vars: EnvVars,
extensions: FrontendHeartbeatExtensions,
cancellation: CancellationToken,
generation: Arc<AtomicU64>,
}
impl HeartbeatTask {
pub fn new(
peer_addr: String,
opts: &FrontendOptions,
meta_client: Arc<MetaClient>,
resp_handler_executor: HeartbeatResponseHandlerExecutorRef,
resource_stat: ResourceStatRef,
) -> Self {
HeartbeatTask {
peer_addr,
meta_client,
resp_handler_executor,
start_time_ms: common_time::util::current_time_millis() as u64,
resource_stat,
env_vars: EnvVars::from_config(&opts.heartbeat_env_vars),
impl HeartbeatRunner {
async fn connect(&self) -> Result<Option<HeartbeatConnection>> {
tokio::select! {
_ = self.cancellation.cancelled() => Ok(None),
connection = self.connector.connect() => connection.map(Some),
}
}
pub async fn start(&self) -> Result<()> {
let (req_sender, resp_stream, config) = self
.meta_client
.heartbeat()
.await
.context(error::CreateMetaHeartbeatStreamSnafu)?;
info!("Heartbeat started with Metasrv config: {}", config);
let (outgoing_tx, outgoing_rx) = mpsc::channel(16);
let mailbox = Arc::new(HeartbeatMailbox::new(outgoing_tx));
self.start_handle_resp_stream(resp_stream, mailbox, config.retry_interval);
self.start_heartbeat_report(req_sender, outgoing_rx, config.interval);
Ok(())
fn next_generation(&self) -> u64 {
self.generation
.fetch_add(1, Ordering::AcqRel)
.wrapping_add(1)
}
fn start_handle_resp_stream(
&self,
mut resp_stream: HeartbeatStream,
mailbox: MailboxRef,
retry_interval: Duration,
) {
let capture_self = self.clone();
async fn notify_connected(&self, generation: u64) -> bool {
for extension in self.extensions.extensions() {
let result = tokio::select! {
_ = self.cancellation.cancelled() => return false,
result = extension.connected(generation) => result,
};
if let Err(error) = result {
error!(error; "Heartbeat extension '{}' failed its connected callback", extension.name());
}
}
true
}
async fn shutdown_extensions(&self) {
for extension in self.extensions.extensions() {
if let Err(error) = extension.shutdown().await {
error!(error; "Failed to shut down heartbeat extension '{}'", extension.name());
}
}
}
async fn run(self, mut connection: HeartbeatConnection) {
loop {
let retry_interval = connection.config.retry_interval;
if self.run_connection(connection).await == ConnectionEnd::Shutdown {
return;
}
let _handle = common_runtime::spawn_hb(async move {
loop {
match resp_stream.message().await {
Ok(Some(resp)) => {
debug!("Receiving heartbeat response: {:?}", resp);
if let Some(message) = &resp.mailbox_message {
info!("Received mailbox message: {message:?}");
}
let ctx = HeartbeatResponseHandlerContext::new(mailbox.clone(), resp);
if let Err(e) = capture_self.handle_response(ctx).await {
error!(e; "Error while handling heartbeat response");
HEARTBEAT_RECV_COUNT
.with_label_values(&["processing_error"])
.inc();
} else {
HEARTBEAT_RECV_COUNT.with_label_values(&["success"]).inc();
}
}
Ok(None) => {
warn!("Heartbeat response stream closed");
capture_self.start_with_retry(retry_interval).await;
break;
}
Err(e) => {
HEARTBEAT_RECV_COUNT.with_label_values(&["error"]).inc();
error!(e; "Occur error while reading heartbeat response");
capture_self.start_with_retry(retry_interval).await;
if !self.wait_retry(retry_interval).await {
return;
}
info!("Try to re-establish the heartbeat connection to metasrv.");
match self.connect().await {
Ok(Some(next)) => {
let generation = self.next_generation();
if !self.notify_connected(generation).await {
return;
}
connection = next;
break;
}
Ok(None) => return,
Err(error) => {
error!(error; "Failed to re-establish heartbeat connection to metasrv");
}
}
}
});
}
}
async fn wait_retry(&self, retry_interval: Duration) -> bool {
tokio::select! {
_ = self.cancellation.cancelled() => false,
_ = tokio::time::sleep(retry_interval) => true,
}
}
async fn run_connection(&self, connection: HeartbeatConnection) -> ConnectionEnd {
let (outgoing_tx, outgoing_rx) = mpsc::channel(16);
let mailbox = Arc::new(HeartbeatMailbox::new(outgoing_tx));
let report =
self.report_heartbeats(connection.sender, outgoing_rx, connection.config.interval);
let responses = self.handle_responses(connection.stream, mailbox);
tokio::pin!(report);
tokio::pin!(responses);
tokio::select! {
_ = self.cancellation.cancelled() => ConnectionEnd::Shutdown,
end = &mut report => end,
end = &mut responses => end,
}
}
async fn handle_responses(
&self,
mut stream: Box<dyn HeartbeatResponseStream>,
mailbox: MailboxRef,
) -> ConnectionEnd {
loop {
let response = tokio::select! {
_ = self.cancellation.cancelled() => return ConnectionEnd::Shutdown,
response = stream.message() => response,
};
match response {
Ok(Some(response)) => {
debug!("Receiving heartbeat response: {:?}", response);
if let Some(message) = &response.mailbox_message {
info!("Received mailbox message: {message:?}");
}
let context = HeartbeatResponseHandlerContext::new(mailbox.clone(), response);
let result = tokio::select! {
_ = self.cancellation.cancelled() => return ConnectionEnd::Shutdown,
result = self.handle_response(context) => result,
};
if let Err(error) = result {
error!(error; "Error while handling heartbeat response");
HEARTBEAT_RECV_COUNT
.with_label_values(&["processing_error"])
.inc();
} else {
HEARTBEAT_RECV_COUNT.with_label_values(&["success"]).inc();
}
}
Ok(None) => {
warn!("Heartbeat response stream closed");
return ConnectionEnd::Reconnect;
}
Err(error) => {
HEARTBEAT_RECV_COUNT.with_label_values(&["error"]).inc();
error!(error; "Occur error while reading heartbeat response");
return ConnectionEnd::Reconnect;
}
}
}
}
async fn report_heartbeats(
&self,
sender: Arc<dyn HeartbeatRequestSender>,
mut outgoing_rx: Receiver<OutgoingMessage>,
report_interval: Duration,
) -> ConnectionEnd {
let total_cpu_millicores = self.resource_stat.get_total_cpu_millicores();
let total_memory_bytes = self.resource_stat.get_total_memory_bytes();
let mut extensions = HashMap::new();
self.env_vars.into_extensions(&mut extensions);
let heartbeat_request = HeartbeatRequest {
peer: Some(Peer {
// Metasrv calculates the frontend id by hashing this reachable address.
id: 0,
addr: self.peer_addr.clone(),
}),
info: Self::build_node_info(
self.start_time_ms,
total_cpu_millicores,
total_memory_bytes,
),
node_workloads: Some(NodeWorkloads::Frontend(FrontendWorkloads { types: vec![] })),
extensions,
..Default::default()
};
let sleep = tokio::time::sleep(Duration::ZERO);
tokio::pin!(sleep);
loop {
let request = tokio::select! {
_ = self.cancellation.cancelled() => return ConnectionEnd::Shutdown,
message = outgoing_rx.recv() => {
if let Some(message) = message {
Self::new_heartbeat_request(&heartbeat_request, Some(message), 0, 0)
} else {
warn!("Sender has been dropped, exiting the heartbeat loop");
return ConnectionEnd::Reconnect;
}
}
_ = &mut sleep => {
sleep.as_mut().reset(Instant::now() + report_interval);
Self::new_heartbeat_request(
&heartbeat_request,
None,
self.resource_stat.get_cpu_usage_millicores(),
self.resource_stat.get_memory_usage_bytes(),
)
}
};
if let Some(mut request) = request {
if !self.add_request_extensions(&mut request).await {
return ConnectionEnd::Shutdown;
}
debug!(
"Sending a heartbeat request to metasrv, content: {:?}",
request
);
let result = tokio::select! {
_ = self.cancellation.cancelled() => return ConnectionEnd::Shutdown,
result = sender.send(request) => result,
};
if let Err(error) = result {
error!(error; "Failed to send heartbeat to metasrv");
return ConnectionEnd::Reconnect;
}
HEARTBEAT_SENT_COUNT.inc();
}
}
}
async fn add_request_extensions(&self, request: &mut HeartbeatRequest) -> bool {
for extension in self.extensions.extensions() {
let generated = tokio::select! {
_ = self.cancellation.cancelled() => return false,
generated = extension.request_extensions() => generated,
};
let generated = match generated {
Ok(generated) => generated,
Err(error) => {
error!(error; "Heartbeat extension '{}' failed to generate request extensions", extension.name());
continue;
}
};
if let Some(key) = generated
.keys()
.find(|key| request.extensions.contains_key(*key))
{
warn!(
"Heartbeat extension '{}' produced conflicting key '{}'; discarding its output",
extension.name(),
key
);
continue;
}
request.extensions.extend(generated);
}
true
}
fn new_heartbeat_request(
@@ -138,8 +500,8 @@ impl HeartbeatTask {
) -> Option<HeartbeatRequest> {
let mailbox_message = match message.map(outgoing_message_to_mailbox_message) {
Some(Ok(message)) => Some(message),
Some(Err(e)) => {
error!(e; "Failed to encode mailbox messages");
Some(Err(error)) => {
error!(error; "Failed to encode mailbox messages");
return None;
}
None => None,
@@ -149,12 +511,10 @@ impl HeartbeatTask {
mailbox_message,
..heartbeat_request.clone()
};
if let Some(info) = heartbeat_request.info.as_mut() {
info.memory_usage_bytes = memory_usage;
info.cpu_usage_millicores = cpu_usage;
}
Some(heartbeat_request)
}
@@ -165,7 +525,6 @@ impl HeartbeatTask {
total_memory_bytes: i64,
) -> Option<NodeInfo> {
let build_info = common_version::build_info();
Some(NodeInfo {
version: build_info.version.to_string(),
git_commit: build_info.commit_short.to_string(),
@@ -184,90 +543,156 @@ impl HeartbeatTask {
})
}
fn start_heartbeat_report(
&self,
req_sender: HeartbeatSender,
mut outgoing_rx: Receiver<OutgoingMessage>,
report_interval: Duration,
) {
let start_time_ms = self.start_time_ms;
let self_peer = Some(Peer {
// The node id will be actually calculated from its address (by hashing the address
// string) in the metasrv. So it can be set to 0 here, as a placeholder.
id: 0,
addr: self.peer_addr.clone(),
});
let total_cpu_millicores = self.resource_stat.get_total_cpu_millicores();
let total_memory_bytes = self.resource_stat.get_total_memory_bytes();
let resource_stat = self.resource_stat.clone();
let env_vars = self.env_vars.clone();
common_runtime::spawn_hb(async move {
let sleep = tokio::time::sleep(Duration::from_millis(0));
tokio::pin!(sleep);
let mut extensions = std::collections::HashMap::new();
env_vars.into_extensions(&mut extensions);
let heartbeat_request = HeartbeatRequest {
peer: self_peer,
info: Self::build_node_info(
start_time_ms,
total_cpu_millicores,
total_memory_bytes,
),
node_workloads: Some(NodeWorkloads::Frontend(FrontendWorkloads { types: vec![] })),
extensions,
..Default::default()
};
loop {
let req = tokio::select! {
message = outgoing_rx.recv() => {
if let Some(message) = message {
Self::new_heartbeat_request(&heartbeat_request, Some(message), 0, 0)
} else {
warn!("Sender has been dropped, exiting the heartbeat loop");
// Receives None that means Sender was dropped, we need to break the current loop
break
}
}
_ = &mut sleep => {
sleep.as_mut().reset(Instant::now() + report_interval);
Self::new_heartbeat_request(&heartbeat_request, None, resource_stat.get_cpu_usage_millicores(), resource_stat.get_memory_usage_bytes())
}
};
if let Some(req) = req {
if let Err(e) = req_sender.send(req.clone()).await {
error!(e; "Failed to send heartbeat to metasrv");
break;
} else {
HEARTBEAT_SENT_COUNT.inc();
debug!("Send a heartbeat request to metasrv, content: {:?}", req);
}
}
}
});
}
async fn handle_response(&self, ctx: HeartbeatResponseHandlerContext) -> Result<()> {
async fn handle_response(&self, context: HeartbeatResponseHandlerContext) -> Result<()> {
self.resp_handler_executor
.handle(ctx)
.handle(context)
.await
.context(error::HandleHeartbeatResponseSnafu)
}
}
async fn start_with_retry(&self, retry_interval: Duration) {
loop {
tokio::time::sleep(retry_interval).await;
#[derive(Debug, PartialEq, Eq)]
enum ConnectionEnd {
Reconnect,
Shutdown,
}
info!("Try to re-establish the heartbeat connection to metasrv.");
/// The frontend task that sends [`HeartbeatRequest`] values to metasrv in the background.
#[derive(Clone)]
pub struct HeartbeatTask {
runner: HeartbeatRunner,
start_lock: Arc<Mutex<()>>,
shutdown_lock: Arc<Mutex<()>>,
supervisor: Arc<Mutex<Option<common_runtime::JoinHandle<()>>>>,
}
if self.start().await.is_ok() {
break;
}
impl HeartbeatTask {
pub fn new(
peer_addr: String,
opts: &FrontendOptions,
meta_client: Arc<MetaClient>,
resp_handler_executor: HeartbeatResponseHandlerExecutorRef,
resource_stat: ResourceStatRef,
) -> Self {
Self::new_with_connector(
peer_addr,
opts,
Arc::new(MetaHeartbeatConnector {
client: meta_client,
}),
resp_handler_executor,
resource_stat,
)
}
fn new_with_connector(
peer_addr: String,
opts: &FrontendOptions,
connector: Arc<dyn HeartbeatConnector>,
resp_handler_executor: HeartbeatResponseHandlerExecutorRef,
resource_stat: ResourceStatRef,
) -> Self {
Self {
runner: HeartbeatRunner {
peer_addr,
connector,
resp_handler_executor,
start_time_ms: common_time::util::current_time_millis() as u64,
resource_stat,
env_vars: EnvVars::from_config(&opts.heartbeat_env_vars),
extensions: FrontendHeartbeatExtensions::default(),
cancellation: CancellationToken::new(),
generation: Arc::new(AtomicU64::new(0)),
},
start_lock: Arc::new(Mutex::new(())),
shutdown_lock: Arc::new(Mutex::new(())),
supervisor: Arc::new(Mutex::new(None)),
}
}
/// Installs the extensions registered before heartbeat startup.
pub fn with_extensions(mut self, extensions: FrontendHeartbeatExtensions) -> Self {
self.runner.extensions = extensions;
self
}
/// Establishes the initial heartbeat connection and starts its background supervisor.
pub async fn start(&self) -> Result<()> {
let _start_guard = self.start_lock.lock().await;
if self.runner.cancellation.is_cancelled() {
return Ok(());
}
let finished = {
let mut supervisor = self.supervisor.lock().await;
match supervisor.as_ref() {
Some(handle) if !handle.is_finished() => return Ok(()),
Some(_) => supervisor.take(),
None => None,
}
};
if let Some(handle) = finished
&& let Err(error) = handle.await
&& !error.is_cancelled()
{
error!(error; "Heartbeat supervisor join failed");
}
let Some(connection) = self.runner.connect().await? else {
return Ok(());
};
info!(
"Heartbeat started with Metasrv config: {}",
connection.config
);
let generation = self.runner.next_generation();
if !self.runner.notify_connected(generation).await {
return Ok(());
}
let runner = self.runner.clone();
let handle = common_runtime::spawn_hb(async move {
runner.run(connection).await;
});
*self.supervisor.lock().await = Some(handle);
Ok(())
}
/// Cancels and joins heartbeat I/O and all registered extension lifecycles.
pub async fn shutdown(&self) {
let _shutdown_guard = self.shutdown_lock.lock().await;
if self.runner.cancellation.is_cancelled() {
return;
}
self.runner.cancellation.cancel();
// Wait for a concurrently running handshake or connected callback to observe cancellation.
let _start_guard = self.start_lock.lock().await;
let handle = self.supervisor.lock().await.take();
if let Some(handle) = handle
&& let Err(error) = handle.await
&& !error.is_cancelled()
{
error!(error; "Heartbeat supervisor join failed");
}
self.runner.shutdown_extensions().await;
}
#[cfg(test)]
fn generation(&self) -> u64 {
self.runner.generation.load(Ordering::Acquire)
}
#[cfg(test)]
async fn has_supervisor(&self) -> bool {
self.supervisor.lock().await.is_some()
}
#[cfg(test)]
pub(crate) fn is_shutdown(&self) -> bool {
self.runner.cancellation.is_cancelled()
}
}
pub(crate) fn frontend_peer_addr(opts: &FrontendOptions) -> String {
File diff suppressed because it is too large Load Diff
+1
View File
@@ -405,6 +405,7 @@ impl HeartbeatHandlerGroup {
region_lease: acc.region_lease,
mailbox_message,
heartbeat_config,
extensions: Default::default(),
};
Ok(res)
}
+17
View File
@@ -12,11 +12,14 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use std::sync::Arc;
use auth::{DefaultPermissionChecker, PermissionCheckerRef, UserProviderRef};
use common_base::Plugins;
use common_meta::cache::CacheRegistryBuilder;
use frontend::error::{IllegalAuthConfigSnafu, Result};
use frontend::frontend::FrontendOptions;
use frontend::heartbeat::FrontendHeartbeatExtensions;
use frontend::instance::Instance;
use frontend::instance::builder::FrontendBuilder;
use snafu::ResultExt;
@@ -64,6 +67,20 @@ pub async fn setup_frontend_plugins_post_build(
Ok(())
}
/// Sets up heartbeat extensions after the frontend [`Instance`] is available.
///
/// Implementations may use the instance to construct extensions and register them in the
/// [`FrontendHeartbeatExtensions`] stored in `plugins`. Registrations must be idempotent because
/// plugin setup can be invoked more than once.
pub async fn setup_frontend_heartbeat_extensions(
plugins: &mut Plugins,
_plugin_options: &[PluginOptions],
_instance: &Arc<Instance>,
) -> Result<()> {
plugins.get_or_insert(FrontendHeartbeatExtensions::default);
Ok(())
}
pub async fn start_frontend_plugins(_instance: &Instance) -> Result<()> {
Ok(())
}
+2 -1
View File
@@ -28,7 +28,8 @@ pub use flownode::{
setup_flownode_plugins_post_build, setup_flownode_plugins_pre_build, start_flownode_plugins,
};
pub use frontend::{
setup_frontend_plugins_post_build, setup_frontend_plugins_pre_build, start_frontend_plugins,
setup_frontend_heartbeat_extensions, setup_frontend_plugins_post_build,
setup_frontend_plugins_pre_build, start_frontend_plugins,
};
pub use meta_srv::{
setup_metasrv_plugins_post_build, setup_metasrv_plugins_pre_build, start_metasrv_plugins,