use super::metrics::ServerMetrics;
use super::router;
use super::tls::{self, TlsConfigError, TlsListener};
use crate::config::{ServerConfig, ServerConfigError};
use axum::Router;
use loonfs::metrics::{JsonlObjectStoreMetricsRecorder, ObjectStoreMetricsRecorder};
use loonfs::{
FsAdmin, FsReader, FsWriter, MaintenanceHandle, MaintenanceJob, MaintenanceProbe,
SharedObjectStore, TraceMode, TraceStoreKind,
};
use loonfs_api::NamespaceId;
use loonfs_grep::{GrepGcJob, GrepMaintenanceJob, GrepService, GrepWorker, GREP_INDEX_JOB};
use loonfs_objectstore::presign::ObjectTransferIssuer;
use std::ffi::OsString;
use std::net::SocketAddr;
use std::sync::Arc;
use thiserror::Error;
use tokio::sync::Semaphore;
const OBJECT_STORE_METRICS_JSONL_ENV: &str = "LOONFS_OBJECT_STORE_METRICS_JSONL";
#[derive(Clone)]
pub(super) struct AppState {
pub(super) config: Arc<ServerConfig>,
pub(super) writer: FsWriter,
pub(super) reader: FsReader,
pub(super) admin: FsAdmin,
pub(super) probe_store: SharedObjectStore,
pub(super) transfer_issuer: Option<Arc<dyn ObjectTransferIssuer>>,
pub(super) grep_worker: Option<GrepWorker<SharedObjectStore>>,
pub(super) grep_service: Option<Arc<GrepService>>,
pub(super) grep_maintenance: Option<GrepMaintenance>,
pub(super) upload_permits: Arc<Semaphore>,
pub(super) download_permits: Arc<Semaphore>,
pub(super) metrics: Arc<ServerMetrics>,
}
#[derive(Clone)]
pub(super) struct GrepMaintenance {
handle: MaintenanceHandle,
job: Arc<GrepMaintenanceJob<SharedObjectStore>>,
}
impl GrepMaintenance {
pub(super) fn nudge(&self, namespace_id: &NamespaceId) {
self.handle.nudge(GREP_INDEX_JOB, namespace_id);
}
pub(super) async fn nudge_if_behind(&self, namespace_id: &NamespaceId) {
if matches!(
self.job.probe(namespace_id).await,
Ok(MaintenanceProbe::Due)
) {
self.nudge(namespace_id);
}
}
}
pub async fn app(config: ServerConfig) -> Result<(Router, FsWriter), ServerConfigError> {
config.validate()?;
let store = config.object_store()?;
let transfer_issuer = config
.store
.direct_put_is_proven()
.then(|| store.transfer_issuer())
.flatten();
let store = store.into_shared();
let (router, state) =
app_with_store_and_transfer_issuer(config, store, transfer_issuer).await?;
Ok((router, state.writer))
}
#[cfg(test)]
pub(super) async fn app_with_store(
config: ServerConfig,
store: SharedObjectStore,
) -> Result<Router, ServerConfigError> {
Ok(app_with_store_and_transfer_issuer(config, store, None)
.await?
.0)
}
#[cfg(test)]
pub(super) async fn app_with_store_and_state(
config: ServerConfig,
store: SharedObjectStore,
) -> Result<(Router, AppState), ServerConfigError> {
app_with_store_and_transfer_issuer(config, store, None).await
}
pub(super) async fn app_with_store_and_transfer_issuer(
config: ServerConfig,
store: SharedObjectStore,
transfer_issuer: Option<Arc<dyn ObjectTransferIssuer>>,
) -> Result<(Router, AppState), ServerConfigError> {
let metrics = ServerMetrics::new();
let maintains_grep_index =
config.maintenance.registers_automatic_jobs() && config.grep.mode.maintains_index();
let (writer, reader, admin) = build_handles(
&config,
store,
&metrics,
std::env::var_os(OBJECT_STORE_METRICS_JSONL_ENV),
)
.await?;
let probe_store = writer.object_store();
let grep_worker = (config.grep.mode.serves_grep() || config.grep.mode.maintains_index())
.then(|| GrepWorker::new(writer.object_store(), reader.clone(), admin.clone()));
let grep_service = config
.grep
.mode
.serves_grep()
.then(|| Arc::new(GrepService::new()));
let grep_maintenance = if maintains_grep_index {
let policy = config
.grep
.worker_config()
.build_policy()
.map_err(|error| ServerConfigError::InvalidField {
field: "grep",
reason: error.to_string(),
})?;
let job = Arc::new(GrepMaintenanceJob::new(
grep_worker
.as_ref()
.expect("an index-maintaining deployment composes a grep worker")
.clone(),
policy,
));
writer
.register_maintenance_job(job.clone())
.map_err(|error| ServerConfigError::InvalidField {
field: "grep",
reason: error.to_string(),
})?;
writer
.register_maintenance_job(Arc::new(GrepGcJob::new(
grep_worker
.as_ref()
.expect("an index-maintaining deployment composes a grep worker")
.clone(),
)))
.map_err(|error| ServerConfigError::InvalidField {
field: "grep",
reason: error.to_string(),
})?;
Some(GrepMaintenance {
handle: writer.maintenance(),
job,
})
} else {
None
};
let config = Arc::new(config);
let state = AppState {
upload_permits: Arc::new(Semaphore::new(
config.max_concurrent_uploads.min(Semaphore::MAX_PERMITS),
)),
download_permits: Arc::new(Semaphore::new(
config.max_concurrent_downloads.min(Semaphore::MAX_PERMITS),
)),
config,
writer,
reader,
admin,
probe_store,
transfer_issuer,
grep_worker,
grep_service,
grep_maintenance,
metrics,
};
Ok((router(state.clone()), state))
}
#[cfg(test)]
pub(super) async fn build_handles_with_metrics_jsonl_path(
config: &ServerConfig,
store: SharedObjectStore,
metrics_jsonl_path: Option<OsString>,
) -> Result<(FsWriter, FsReader, FsAdmin), ServerConfigError> {
build_handles(config, store, &ServerMetrics::new(), metrics_jsonl_path).await
}
async fn build_handles(
config: &ServerConfig,
store: SharedObjectStore,
metrics: &ServerMetrics,
metrics_jsonl_path: Option<OsString>,
) -> Result<(FsWriter, FsReader, FsAdmin), ServerConfigError> {
let trace_store_kind = TraceStoreKind::from(config.store.kind());
let samples = object_store_metrics_recorder(metrics_jsonl_path)?;
let runtime_error = |error: loonfs::RuntimeError| ServerConfigError::InvalidField {
field: "runtime",
reason: error.to_string(),
};
let mut writer_builder = FsWriter::builder_with_store(store.clone())
.writer_id(config.writer_id.clone())
.background_work(config.maintenance.background_work())
.min_publish_interval_ms(config.min_publish_interval_ms)
.max_read_content_bytes(config.max_download_bytes)
.max_concurrent_maintenance(config.max_concurrent_maintenance)
.runtime_cache(config.runtime_cache_config())
.trace_mode(TraceMode::Remote)
.trace_store_kind(trace_store_kind)
.metrics_recorder(metrics.recorder());
if let Some(samples) = &samples {
writer_builder = writer_builder.object_store_metrics_recorder(Arc::clone(samples));
}
let writer = writer_builder.build().await.map_err(runtime_error)?;
let reader = writer.reader();
let mut admin_builder = FsAdmin::builder_with_store(store)
.actor_id(format!("{}-admin", config.writer_id))
.runtime_cache(config.runtime_cache_config())
.shared_metadata_table_cache(&writer)
.trace_mode(TraceMode::Remote)
.trace_store_kind(trace_store_kind)
.metrics_recorder(metrics.recorder());
if let Some(samples) = samples {
admin_builder = admin_builder.object_store_metrics_recorder(samples);
}
let admin = admin_builder.build().await.map_err(runtime_error)?;
Ok((writer, reader, admin))
}
fn object_store_metrics_recorder(
metrics_jsonl_path: Option<OsString>,
) -> Result<Option<Arc<dyn ObjectStoreMetricsRecorder>>, ServerConfigError> {
let Some(path) = metrics_jsonl_path else {
return Ok(None);
};
if path.is_empty() {
return Ok(None);
}
let path = std::path::PathBuf::from(path);
JsonlObjectStoreMetricsRecorder::create(&path)
.map(|recorder| Some(Arc::new(recorder) as Arc<dyn ObjectStoreMetricsRecorder>))
.map_err(|error| ServerConfigError::InvalidField {
field: OBJECT_STORE_METRICS_JSONL_ENV,
reason: error.to_string(),
})
}
#[derive(Debug, Error)]
pub enum ServeError {
#[error("invalid server config: {0}")]
Config(#[from] ServerConfigError),
#[error("failed to bind `{addr}`: {source}")]
Bind {
addr: SocketAddr,
#[source]
source: std::io::Error,
},
#[error("failed to load the configured TLS identity: {0}")]
Tls(#[source] TlsConfigError),
#[error("server failed while serving requests: {0}")]
Serve(#[source] std::io::Error),
#[error("background work did not settle during shutdown: {0}")]
Shutdown(#[source] loonfs::RuntimeError),
}
pub async fn serve(config: ServerConfig) -> Result<(), ServeError> {
serve_with_shutdown(config, shutdown_signal()).await
}
pub async fn serve_with_shutdown(
config: ServerConfig,
shutdown: impl std::future::Future<Output = ()> + Send + 'static,
) -> Result<(), ServeError> {
let bind = config.bind_addr()?;
let tls = config
.tls
.as_ref()
.map(tls::server_config)
.transpose()
.map_err(ServeError::Tls)?;
let listener = tokio::net::TcpListener::bind(bind)
.await
.map_err(|source| ServeError::Bind { addr: bind, source })?;
match tls {
Some(tls) => serve_on(TlsListener::new(listener, tls), config, shutdown).await,
None => serve_on(listener, config, shutdown).await,
}
}
pub(super) async fn serve_on<L>(
listener: L,
config: ServerConfig,
shutdown: impl std::future::Future<Output = ()> + Send + 'static,
) -> Result<(), ServeError>
where
L: axum::serve::Listener<Addr = SocketAddr>,
{
let (router, writer) = app(config).await?;
axum::serve(listener, router)
.with_graceful_shutdown(shutdown)
.await
.map_err(ServeError::Serve)?;
writer.shutdown().await.map_err(ServeError::Shutdown)
}
async fn shutdown_signal() {
let ctrl_c = async {
tokio::signal::ctrl_c()
.await
.expect("ctrl-c handler should install");
};
#[cfg(unix)]
let terminate = async {
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
.expect("SIGTERM handler should install")
.recv()
.await;
};
#[cfg(not(unix))]
let terminate = std::future::pending::<()>();
tokio::select! {
() = ctrl_c => {}
_ = terminate => {}
}
}