use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::Mutex;
use tracing::{debug, warn};
use crate::block::router::WorkerRouter;
use crate::client::metrics_master::MetricsClient;
use crate::client::metrics_master::MetricsMasterClient;
use crate::client::{
create_master_inquire_client, MasterClient, MasterInquireClient, WorkerClientPool,
WorkerManagerClient,
};
use crate::config::{ConfigRefresher, GoosefsConfig, TransparentAccelerationSwitch};
use crate::error::{Error, Result};
use crate::metrics::heartbeat::{resolve_app_id, HeartbeatTask};
use crate::metrics::reporter::ClientMetricsReporter;
const REFRESH_CHECK_INTERVAL: Duration = Duration::from_secs(30);
const CONFIG_REFRESH_INTERVAL: Duration = Duration::from_secs(60);
pub struct FileSystemContext {
config: Arc<GoosefsConfig>,
master: Arc<MasterClient>,
worker_manager: Arc<WorkerManagerClient>,
worker_pool: Arc<WorkerClientPool>,
worker_router: Arc<WorkerRouter>,
inquire_client: Arc<dyn MasterInquireClient>,
config_refresher: Arc<ConfigRefresher>,
closed: Arc<AtomicBool>,
worker_refresh_task: Mutex<Option<tokio::task::JoinHandle<()>>>,
config_refresh_task: Mutex<Option<tokio::task::JoinHandle<()>>>,
metrics_heartbeat: Mutex<Option<Arc<HeartbeatTask>>>,
}
impl FileSystemContext {
pub async fn connect(config: GoosefsConfig) -> Result<Arc<Self>> {
let config = Arc::new(config);
let inquire_client = create_master_inquire_client(&config);
let (master_res, wm_res) = tokio::join!(
MasterClient::connect_with_inquire(&config, inquire_client.clone()),
WorkerManagerClient::connect_with_inquire(&config, inquire_client.clone()),
);
let master = Arc::new(master_res?);
let worker_manager = Arc::new(wm_res?);
let workers = worker_manager.get_worker_info_list().await?;
if workers.is_empty() {
return Err(Error::NoWorkerAvailable {
message: "no workers available at startup".to_string(),
});
}
debug!(count = workers.len(), "initial worker list fetched");
let worker_router = Arc::new(WorkerRouter::with_ttls(
Duration::from_secs(60), Duration::from_secs(30), ));
worker_router.update_workers(workers).await;
let worker_pool = WorkerClientPool::new_shared((*config).clone());
let ctx = Arc::new(Self {
config: config.clone(),
master,
worker_manager,
worker_pool,
worker_router,
inquire_client,
config_refresher: Arc::new(ConfigRefresher::from_config(&config)),
closed: Arc::new(AtomicBool::new(false)),
worker_refresh_task: Mutex::new(None),
config_refresh_task: Mutex::new(None),
metrics_heartbeat: Mutex::new(None),
});
ctx.clone().start_worker_refresh_task().await;
ctx.clone().start_config_refresh_task().await;
ctx.clone().start_metrics_heartbeat_task().await?;
Ok(ctx)
}
pub fn acquire_master(&self) -> Arc<MasterClient> {
self.master.clone()
}
pub fn acquire_worker_manager(&self) -> Arc<WorkerManagerClient> {
self.worker_manager.clone()
}
pub fn acquire_worker_pool(&self) -> Arc<WorkerClientPool> {
self.worker_pool.clone()
}
pub fn acquire_router(&self) -> Arc<WorkerRouter> {
self.worker_router.clone()
}
pub fn acquire_inquire_client(&self) -> Arc<dyn MasterInquireClient> {
self.inquire_client.clone()
}
pub fn config(&self) -> &GoosefsConfig {
&self.config
}
pub fn acquire_config_refresher(&self) -> Arc<ConfigRefresher> {
self.config_refresher.clone()
}
pub fn refresh_transparent_acceleration_switch(&self) -> TransparentAccelerationSwitch {
self.config_refresher
.refresh_transparent_acceleration_switch()
}
pub async fn close(&self) -> Result<()> {
if self.closed.swap(true, Ordering::SeqCst) {
return Ok(()); }
let worker_handle = self.worker_refresh_task.lock().await.take();
if let Some(h) = worker_handle {
h.abort();
debug!("worker refresh task aborted");
}
let config_handle = self.config_refresh_task.lock().await.take();
if let Some(h) = config_handle {
h.abort();
debug!("config refresh task aborted");
}
if let Some(task) = self.metrics_heartbeat.lock().await.take() {
task.shutdown().await;
debug!("metrics heartbeat task shut down");
}
Ok(())
}
pub fn is_closed(&self) -> bool {
self.closed.load(Ordering::SeqCst)
}
async fn start_worker_refresh_task(self: Arc<Self>) {
let worker_router = self.worker_router.clone();
let worker_manager = self.worker_manager.clone();
let closed = self.closed.clone();
let handle = tokio::spawn(async move {
loop {
tokio::time::sleep(REFRESH_CHECK_INTERVAL).await;
if closed.load(Ordering::SeqCst) {
debug!("worker refresh task: context closed, exiting");
break;
}
if worker_router.needs_refresh().await {
if let Err(e) = worker_router.refresh_workers(&worker_manager).await {
warn!("worker refresh failed: {}", e);
} else {
debug!("worker list refreshed by background task");
}
}
}
});
*self.worker_refresh_task.lock().await = Some(handle);
}
async fn start_config_refresh_task(self: Arc<Self>) {
let config_refresher = self.config_refresher.clone();
let closed = self.closed.clone();
let handle = tokio::spawn(async move {
let switch = config_refresher.refresh_transparent_acceleration_switch();
debug!(
transparent_acceleration_enabled = switch.enabled,
cosranger_enabled = switch.cosranger_enabled,
"config refresh: initial load completed"
);
loop {
tokio::time::sleep(CONFIG_REFRESH_INTERVAL).await;
if closed.load(Ordering::SeqCst) {
debug!("config refresh task: context closed, exiting");
break;
}
let switch = config_refresher.refresh_transparent_acceleration_switch();
debug!(
transparent_acceleration_enabled = switch.enabled,
cosranger_enabled = switch.cosranger_enabled,
"config refresh check completed"
);
}
});
*self.config_refresh_task.lock().await = Some(handle);
}
}
impl Drop for FileSystemContext {
fn drop(&mut self) {
self.closed.store(true, Ordering::SeqCst);
if let Ok(mut guard) = self.worker_refresh_task.try_lock() {
if let Some(h) = guard.take() {
h.abort();
}
}
if let Ok(mut guard) = self.config_refresh_task.try_lock() {
if let Some(h) = guard.take() {
h.abort();
}
}
if let Ok(mut guard) = self.metrics_heartbeat.try_lock() {
guard.take(); }
}
}
impl FileSystemContext {
async fn start_metrics_heartbeat_task(self: Arc<Self>) -> Result<()> {
if !self.config.metrics_enabled {
debug!("metrics disabled — heartbeat task not started");
return Ok(());
}
let mm_client =
MetricsMasterClient::connect_with_inquire(&self.config, self.inquire_client.clone())
.await?;
let reporter = Arc::new(ClientMetricsReporter::default());
let app_id = resolve_app_id(&self.config);
debug!(
app_id = %app_id,
interval_ms = self.config.metrics_heartbeat_interval.as_millis(),
timeout_ms = self.config.metrics_heartbeat_timeout.as_millis(),
"starting metrics heartbeat task"
);
let task = Arc::new(HeartbeatTask::spawn(
Arc::new(mm_client) as Arc<dyn MetricsClient>,
reporter,
app_id,
self.config.metrics_heartbeat_interval,
self.config.metrics_heartbeat_timeout,
self.closed.clone(),
));
*self.metrics_heartbeat.lock().await = Some(task);
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
#[test]
fn test_context_closed_starts_false() {
let closed = Arc::new(AtomicBool::new(false));
assert!(!closed.load(Ordering::SeqCst));
}
#[test]
fn test_context_close_is_idempotent() {
let closed = Arc::new(AtomicBool::new(false));
let was_open = !closed.swap(true, Ordering::SeqCst);
assert!(was_open);
let was_open2 = !closed.swap(true, Ordering::SeqCst);
assert!(!was_open2);
}
#[test]
fn test_refresh_check_interval() {
assert_eq!(REFRESH_CHECK_INTERVAL, Duration::from_secs(30));
}
#[test]
fn test_config_refresh_interval() {
assert_eq!(CONFIG_REFRESH_INTERVAL, Duration::from_secs(60));
}
#[test]
fn test_worker_router_ttls_accepted() {
let router = WorkerRouter::with_ttls(Duration::from_secs(60), Duration::from_secs(30));
drop(router);
}
#[test]
fn test_resolve_app_id_non_empty() {
let config = GoosefsConfig::new("127.0.0.1:9200");
let id = resolve_app_id(&config);
assert!(!id.is_empty());
}
#[test]
fn test_metrics_enabled_flag_is_accessible() {
let config_on = GoosefsConfig::new("127.0.0.1:9200").with_metrics_enabled(true);
assert!(config_on.metrics_enabled);
let config_off = GoosefsConfig::new("127.0.0.1:9200").with_metrics_enabled(false);
assert!(!config_off.metrics_enabled);
}
#[tokio::test]
async fn disabled_no_task_spawn() {
let config = Arc::new(GoosefsConfig::new("127.0.0.1:9200").with_metrics_enabled(false));
assert!(!config.metrics_enabled, "metrics_enabled must be false");
let metrics_heartbeat: Mutex<Option<Arc<HeartbeatTask>>> = Mutex::new(None);
let task_was_spawned = if config.metrics_enabled {
true
} else {
false
};
assert!(
!task_was_spawned,
"metrics_enabled=false must prevent task from being spawned"
);
let guard = metrics_heartbeat.lock().await;
assert!(
guard.is_none(),
"metrics_heartbeat field must remain None when metrics are disabled"
);
}
#[test]
fn metrics_disabled_by_default() {
let config = GoosefsConfig::new("127.0.0.1:9200");
let config_off = GoosefsConfig::new("127.0.0.1:9200").with_metrics_enabled(false);
assert!(
!config_off.metrics_enabled,
"with_metrics_enabled(false) must disable metrics"
);
assert!(
config.metrics_enabled,
"metrics_enabled defaults to true (opt-in enabled by default per Java SDK alignment)"
);
}
}