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::block::short_circuit::{ShortCircuitConfig, ShortCircuitFactory};
use crate::cache::{CacheManager, LocalCacheManager};
use crate::client::metrics_master::MetricsClient;
use crate::client::metrics_master::MetricsMasterClient;
use crate::client::{
create_master_inquire_client, MasterClient, MasterClientPool, MasterInquireClient,
WorkerClientPool, WorkerManagerClient,
};
use crate::config::{ConfigRefresher, GoosefsConfig, TransparentAccelerationSwitch};
use crate::error::Result;
use crate::file_info_cache::FileInfoCache;
use crate::metrics::heartbeat::{resolve_app_id, HeartbeatTask};
#[cfg(feature = "metrics-pushgateway")]
use crate::metrics::pushgateway::{PushgatewayConfig, PushgatewayTask};
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_pool: Arc<MasterClientPool>,
worker_manager: Option<Arc<WorkerManagerClient>>,
worker_pool: Arc<WorkerClientPool>,
worker_router: Arc<WorkerRouter>,
short_circuit: Option<Arc<ShortCircuitFactory>>,
cache_manager: Option<Arc<dyn CacheManager>>,
file_info_cache: Option<Arc<FileInfoCache>>,
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>>>,
#[cfg(feature = "metrics-pushgateway")]
pushgateway_task: Mutex<Option<PushgatewayTask>>,
}
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 (pool_res, wm_res) = tokio::join!(
MasterClientPool::connect_with_inquire(&config, inquire_client.clone()),
WorkerManagerClient::connect_with_inquire(&config, inquire_client.clone()),
);
let master_pool = Arc::new(pool_res?);
let worker_router = Arc::new(WorkerRouter::with_ttls(
Duration::from_secs(60), Duration::from_secs(30), ));
let worker_manager = match wm_res {
Ok(wm) => {
match wm.get_worker_info_list().await {
Ok(workers) => {
if workers.is_empty() {
warn!("WorkerManager returned empty worker list — proceeding without worker discovery");
None
} else {
debug!(count = workers.len(), "initial worker list fetched");
worker_router.update_workers(workers).await;
Some(Arc::new(wm))
}
}
Err(e) => {
warn!("GetWorkerInfoList failed ({}), proceeding without worker discovery. \
Master-only operations (CreateFile, GetStatus, etc.) will still work.", e);
None
}
}
}
Err(e) => {
warn!(
"WorkerManager connection failed ({}), proceeding without worker discovery. \
Master-only operations (CreateFile, GetStatus, etc.) will still work.",
e
);
None
}
};
let worker_pool = WorkerClientPool::new_shared((*config).clone());
let sc_cfg = ShortCircuitConfig::from_config(&config);
let short_circuit = if sc_cfg.enabled {
Some(Arc::new(ShortCircuitFactory::new(
worker_pool.clone(),
worker_router.clone(),
sc_cfg,
)))
} else {
None
};
let cache_manager: Option<Arc<dyn CacheManager>> = if config.client_cache_enabled {
match LocalCacheManager::from_config(&config).await {
Ok(mgr) => {
debug!(
page_size = config.client_cache_page_size,
dirs = ?config.client_cache_dirs,
"client local page cache enabled"
);
Some(mgr as Arc<dyn CacheManager>)
}
Err(e) => {
warn!(error = %e, "failed to init client page cache; continuing without cache");
None
}
}
} else {
None
};
let file_info_cache =
FileInfoCache::maybe_new(config.file_info_cache_ttl, config.file_info_cache_capacity);
if let Some(c) = &file_info_cache {
debug!(
ttl_ms = config.file_info_cache_ttl.as_millis(),
capacity = config.file_info_cache_capacity,
"FileInfo metadata cache enabled (opt-in, §A3), ttl={:?}",
c.ttl(),
);
}
let ctx = Arc::new(Self {
config: config.clone(),
master_pool,
worker_manager,
worker_pool,
worker_router,
short_circuit,
cache_manager,
file_info_cache,
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),
#[cfg(feature = "metrics-pushgateway")]
pushgateway_task: Mutex::new(None),
});
ctx.clone().start_worker_refresh_task().await;
ctx.clone().start_config_refresh_task().await;
ctx.clone().start_metrics_heartbeat_task().await?;
#[cfg(feature = "metrics-pushgateway")]
ctx.clone().start_pushgateway_task().await;
Ok(ctx)
}
pub fn acquire_master(&self) -> Arc<MasterClient> {
self.master_pool.pick()
}
pub fn acquire_master_pool(&self) -> Arc<MasterClientPool> {
self.master_pool.clone()
}
pub fn acquire_worker_manager(&self) -> Option<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_short_circuit(&self) -> Option<Arc<ShortCircuitFactory>> {
self.short_circuit.clone()
}
pub fn acquire_cache_manager(&self) -> Option<Arc<dyn CacheManager>> {
self.cache_manager.clone()
}
pub fn acquire_file_info_cache(&self) -> Option<Arc<FileInfoCache>> {
self.file_info_cache.clone()
}
pub fn invalidate_file_info(&self, path: &str) {
if let Some(cache) = &self.file_info_cache {
cache.invalidate(path);
}
}
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");
}
#[cfg(feature = "metrics-pushgateway")]
if let Some(task) = self.pushgateway_task.lock().await.take() {
task.shutdown().await;
debug!("pushgateway 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 Some(wm) = worker_manager else {
debug!("worker refresh task skipped: no WorkerManager available");
return;
};
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(&wm).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(feature = "metrics-pushgateway")]
async fn start_pushgateway_task(self: Arc<Self>) {
let effective_config = if self.config.pushgateway_enabled {
None
} else {
match GoosefsConfig::from_properties_auto() {
Ok(file_cfg) if file_cfg.pushgateway_enabled => {
debug!(
"pushgateway not enabled in initial config, \
but enabled in properties file — using file config"
);
Some(file_cfg)
}
_ => {
debug!("pushgateway disabled — push task not started");
return;
}
}
};
let cfg = effective_config.as_ref().unwrap_or(&self.config);
let mut pg_config = PushgatewayConfig::new(
cfg.pushgateway_endpoint.clone(),
cfg.pushgateway_job.clone(),
)
.with_push_interval(cfg.pushgateway_push_interval);
if let Some(ref instance) = cfg.pushgateway_instance {
pg_config = pg_config.with_instance(instance.clone());
} else {
let pid = std::process::id();
let ip = Self::resolve_local_ip();
let auto_instance = format!("{}:{}", ip, pid);
debug!(auto_instance = %auto_instance, "auto-generated pushgateway instance");
pg_config = pg_config.with_instance(auto_instance);
}
debug!(
endpoint = %cfg.pushgateway_endpoint,
job = %cfg.pushgateway_job,
interval_ms = cfg.pushgateway_push_interval.as_millis(),
"starting pushgateway push task"
);
let task = PushgatewayTask::spawn(pg_config);
*self.pushgateway_task.lock().await = Some(task);
}
#[cfg(feature = "metrics-pushgateway")]
fn resolve_local_ip() -> String {
use std::net::UdpSocket;
match UdpSocket::bind("0.0.0.0:0") {
Ok(socket) => {
if socket.connect("8.8.8.8:80").is_ok() {
if let Ok(addr) = socket.local_addr() {
return addr.ip().to_string();
}
}
"127.0.0.1".to_string()
}
Err(_) => "127.0.0.1".to_string(),
}
}
}
#[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)"
);
}
#[test]
fn file_info_cache_enabled_by_default() {
let cfg = GoosefsConfig::default();
assert_eq!(
cfg.file_info_cache_ttl,
Duration::from_secs(30),
"default TTL must be 30 s (enabled by default per §A3)"
);
assert!(
crate::file_info_cache::FileInfoCache::maybe_new(
cfg.file_info_cache_ttl,
cfg.file_info_cache_capacity,
)
.is_some(),
"FileInfoCache::maybe_new must return Some when default TTL > 0"
);
}
#[test]
fn file_info_cache_opt_in_produces_live_cache() {
let cfg = GoosefsConfig::new("127.0.0.1:9200")
.with_file_info_cache_ttl(Duration::from_secs(30))
.with_file_info_cache_capacity(256);
assert_eq!(cfg.file_info_cache_ttl, Duration::from_secs(30));
assert_eq!(cfg.file_info_cache_capacity, 256);
let cache = crate::file_info_cache::FileInfoCache::maybe_new(
cfg.file_info_cache_ttl,
cfg.file_info_cache_capacity,
)
.expect("opt-in TTL > 0 must produce a live cache");
assert_eq!(cache.ttl(), Duration::from_secs(30));
}
#[test]
fn file_info_cache_capacity_clamped_to_one() {
let cfg = GoosefsConfig::new("127.0.0.1:9200").with_file_info_cache_capacity(0);
assert_eq!(
cfg.file_info_cache_capacity, 1,
"with_file_info_cache_capacity(0) must clamp to 1"
);
}
}