1use super::metrics::ServerMetrics;
5use super::router;
6use super::tls::{self, TlsConfigError, TlsListener};
7use crate::config::{ServerConfig, ServerConfigError};
8use axum::Router;
9use loonfs::metrics::{JsonlObjectStoreMetricsRecorder, ObjectStoreMetricsRecorder};
10use loonfs::{
11 FsAdmin, FsReader, FsWriter, MaintenanceHandle, MaintenanceJob, MaintenanceProbe,
12 SharedObjectStore, TraceMode, TraceStoreKind,
13};
14use loonfs_api::NamespaceId;
15use loonfs_grep::{GrepGcJob, GrepMaintenanceJob, GrepService, GrepWorker, GREP_INDEX_JOB};
16use loonfs_objectstore::presign::ObjectTransferIssuer;
17use std::ffi::OsString;
18use std::net::SocketAddr;
19use std::sync::Arc;
20use thiserror::Error;
21use tokio::sync::Semaphore;
22
23const OBJECT_STORE_METRICS_JSONL_ENV: &str = "LOONFS_OBJECT_STORE_METRICS_JSONL";
24
25#[derive(Clone)]
34pub(super) struct AppState {
35 pub(super) config: Arc<ServerConfig>,
36 pub(super) writer: FsWriter,
37 pub(super) reader: FsReader,
38 pub(super) admin: FsAdmin,
39 pub(super) probe_store: SharedObjectStore,
44 pub(super) transfer_issuer: Option<Arc<dyn ObjectTransferIssuer>>,
45 pub(super) grep_worker: Option<GrepWorker<SharedObjectStore>>,
46 pub(super) grep_service: Option<Arc<GrepService>>,
50 pub(super) grep_maintenance: Option<GrepMaintenance>,
55 pub(super) upload_permits: Arc<Semaphore>,
60 pub(super) download_permits: Arc<Semaphore>,
64 pub(super) metrics: Arc<ServerMetrics>,
70}
71
72#[derive(Clone)]
78pub(super) struct GrepMaintenance {
79 handle: MaintenanceHandle,
80 job: Arc<GrepMaintenanceJob<SharedObjectStore>>,
81}
82
83impl GrepMaintenance {
84 pub(super) fn nudge(&self, namespace_id: &NamespaceId) {
87 self.handle.nudge(GREP_INDEX_JOB, namespace_id);
88 }
89
90 pub(super) async fn nudge_if_behind(&self, namespace_id: &NamespaceId) {
97 if matches!(
98 self.job.probe(namespace_id).await,
99 Ok(MaintenanceProbe::Due)
100 ) {
101 self.nudge(namespace_id);
102 }
103 }
104}
105
106pub async fn app(config: ServerConfig) -> Result<(Router, FsWriter), ServerConfigError> {
119 config.validate()?;
123 let store = config.object_store()?;
124 let transfer_issuer = config
130 .store
131 .direct_put_is_proven()
132 .then(|| store.transfer_issuer())
133 .flatten();
134 let store = store.into_shared();
135 let (router, state) =
136 app_with_store_and_transfer_issuer(config, store, transfer_issuer).await?;
137 Ok((router, state.writer))
138}
139
140#[cfg(test)]
141pub(super) async fn app_with_store(
142 config: ServerConfig,
143 store: SharedObjectStore,
144) -> Result<Router, ServerConfigError> {
145 Ok(app_with_store_and_transfer_issuer(config, store, None)
146 .await?
147 .0)
148}
149
150#[cfg(test)]
153pub(super) async fn app_with_store_and_state(
154 config: ServerConfig,
155 store: SharedObjectStore,
156) -> Result<(Router, AppState), ServerConfigError> {
157 app_with_store_and_transfer_issuer(config, store, None).await
158}
159
160pub(super) async fn app_with_store_and_transfer_issuer(
161 config: ServerConfig,
162 store: SharedObjectStore,
163 transfer_issuer: Option<Arc<dyn ObjectTransferIssuer>>,
164) -> Result<(Router, AppState), ServerConfigError> {
165 let metrics = ServerMetrics::new();
166 let maintains_grep_index =
170 config.maintenance.registers_automatic_jobs() && config.grep.mode.maintains_index();
171 let (writer, reader, admin) = build_handles(
177 &config,
178 store,
179 &metrics,
180 std::env::var_os(OBJECT_STORE_METRICS_JSONL_ENV),
181 )
182 .await?;
183 let probe_store = writer.object_store();
184 let grep_worker = (config.grep.mode.serves_grep() || config.grep.mode.maintains_index())
189 .then(|| GrepWorker::new(writer.object_store(), reader.clone(), admin.clone()));
190 let grep_service = config
191 .grep
192 .mode
193 .serves_grep()
194 .then(|| Arc::new(GrepService::new()));
195 let grep_maintenance = if maintains_grep_index {
196 let policy = config
197 .grep
198 .worker_config()
199 .build_policy()
200 .map_err(|error| ServerConfigError::InvalidField {
201 field: "grep",
202 reason: error.to_string(),
203 })?;
204 let job = Arc::new(GrepMaintenanceJob::new(
205 grep_worker
206 .as_ref()
207 .expect("an index-maintaining deployment composes a grep worker")
208 .clone(),
209 policy,
210 ));
211 writer
212 .register_maintenance_job(job.clone())
213 .map_err(|error| ServerConfigError::InvalidField {
214 field: "grep",
215 reason: error.to_string(),
216 })?;
217 writer
221 .register_maintenance_job(Arc::new(GrepGcJob::new(
222 grep_worker
223 .as_ref()
224 .expect("an index-maintaining deployment composes a grep worker")
225 .clone(),
226 )))
227 .map_err(|error| ServerConfigError::InvalidField {
228 field: "grep",
229 reason: error.to_string(),
230 })?;
231 Some(GrepMaintenance {
232 handle: writer.maintenance(),
233 job,
234 })
235 } else {
236 None
237 };
238 let config = Arc::new(config);
239 let state = AppState {
240 upload_permits: Arc::new(Semaphore::new(
241 config.max_concurrent_uploads.min(Semaphore::MAX_PERMITS),
242 )),
243 download_permits: Arc::new(Semaphore::new(
244 config.max_concurrent_downloads.min(Semaphore::MAX_PERMITS),
245 )),
246 config,
247 writer,
248 reader,
249 admin,
250 probe_store,
251 transfer_issuer,
252 grep_worker,
253 grep_service,
254 grep_maintenance,
255 metrics,
256 };
257 Ok((router(state.clone()), state))
258}
259
260#[cfg(test)]
261pub(super) async fn build_handles_with_metrics_jsonl_path(
262 config: &ServerConfig,
263 store: SharedObjectStore,
264 metrics_jsonl_path: Option<OsString>,
265) -> Result<(FsWriter, FsReader, FsAdmin), ServerConfigError> {
266 build_handles(config, store, &ServerMetrics::new(), metrics_jsonl_path).await
267}
268
269async fn build_handles(
277 config: &ServerConfig,
278 store: SharedObjectStore,
279 metrics: &ServerMetrics,
280 metrics_jsonl_path: Option<OsString>,
281) -> Result<(FsWriter, FsReader, FsAdmin), ServerConfigError> {
282 let trace_store_kind = TraceStoreKind::from(config.store.kind());
283 let samples = object_store_metrics_recorder(metrics_jsonl_path)?;
284 let runtime_error = |error: loonfs::RuntimeError| ServerConfigError::InvalidField {
285 field: "runtime",
286 reason: error.to_string(),
287 };
288
289 let mut writer_builder = FsWriter::builder_with_store(store.clone())
290 .writer_id(config.writer_id.clone())
291 .background_work(config.maintenance.background_work())
292 .min_publish_interval_ms(config.min_publish_interval_ms)
293 .max_read_content_bytes(config.max_download_bytes)
296 .max_concurrent_maintenance(config.max_concurrent_maintenance)
297 .runtime_cache(config.runtime_cache_config())
298 .trace_mode(TraceMode::Remote)
299 .trace_store_kind(trace_store_kind)
300 .metrics_recorder(metrics.recorder());
301 if let Some(samples) = &samples {
302 writer_builder = writer_builder.object_store_metrics_recorder(Arc::clone(samples));
303 }
304 let writer = writer_builder.build().await.map_err(runtime_error)?;
305 let reader = writer.reader();
306
307 let mut admin_builder = FsAdmin::builder_with_store(store)
308 .actor_id(format!("{}-admin", config.writer_id))
309 .runtime_cache(config.runtime_cache_config())
314 .shared_metadata_table_cache(&writer)
315 .trace_mode(TraceMode::Remote)
316 .trace_store_kind(trace_store_kind)
317 .metrics_recorder(metrics.recorder());
318 if let Some(samples) = samples {
319 admin_builder = admin_builder.object_store_metrics_recorder(samples);
320 }
321 let admin = admin_builder.build().await.map_err(runtime_error)?;
322
323 Ok((writer, reader, admin))
324}
325
326fn object_store_metrics_recorder(
327 metrics_jsonl_path: Option<OsString>,
328) -> Result<Option<Arc<dyn ObjectStoreMetricsRecorder>>, ServerConfigError> {
329 let Some(path) = metrics_jsonl_path else {
330 return Ok(None);
331 };
332 if path.is_empty() {
333 return Ok(None);
334 }
335 let path = std::path::PathBuf::from(path);
336 JsonlObjectStoreMetricsRecorder::create(&path)
337 .map(|recorder| Some(Arc::new(recorder) as Arc<dyn ObjectStoreMetricsRecorder>))
338 .map_err(|error| ServerConfigError::InvalidField {
339 field: OBJECT_STORE_METRICS_JSONL_ENV,
340 reason: error.to_string(),
341 })
342}
343
344#[derive(Debug, Error)]
346pub enum ServeError {
347 #[error("invalid server config: {0}")]
348 Config(#[from] ServerConfigError),
349 #[error("failed to bind `{addr}`: {source}")]
350 Bind {
351 addr: SocketAddr,
352 #[source]
353 source: std::io::Error,
354 },
355 #[error("failed to load the configured TLS identity: {0}")]
356 Tls(#[source] TlsConfigError),
357 #[error("server failed while serving requests: {0}")]
358 Serve(#[source] std::io::Error),
359 #[error("background work did not settle during shutdown: {0}")]
360 Shutdown(#[source] loonfs::RuntimeError),
361}
362
363pub async fn serve(config: ServerConfig) -> Result<(), ServeError> {
368 serve_with_shutdown(config, shutdown_signal()).await
369}
370
371pub async fn serve_with_shutdown(
374 config: ServerConfig,
375 shutdown: impl std::future::Future<Output = ()> + Send + 'static,
376) -> Result<(), ServeError> {
377 let bind = config.bind_addr()?;
378 let tls = config
381 .tls
382 .as_ref()
383 .map(tls::server_config)
384 .transpose()
385 .map_err(ServeError::Tls)?;
386 let listener = tokio::net::TcpListener::bind(bind)
387 .await
388 .map_err(|source| ServeError::Bind { addr: bind, source })?;
389 match tls {
390 Some(tls) => serve_on(TlsListener::new(listener, tls), config, shutdown).await,
391 None => serve_on(listener, config, shutdown).await,
392 }
393}
394
395pub(super) async fn serve_on<L>(
400 listener: L,
401 config: ServerConfig,
402 shutdown: impl std::future::Future<Output = ()> + Send + 'static,
403) -> Result<(), ServeError>
404where
405 L: axum::serve::Listener<Addr = SocketAddr>,
406{
407 let (router, writer) = app(config).await?;
408 axum::serve(listener, router)
409 .with_graceful_shutdown(shutdown)
410 .await
411 .map_err(ServeError::Serve)?;
412 writer.shutdown().await.map_err(ServeError::Shutdown)
418}
419
420async fn shutdown_signal() {
423 let ctrl_c = async {
424 tokio::signal::ctrl_c()
425 .await
426 .expect("ctrl-c handler should install");
427 };
428 #[cfg(unix)]
429 let terminate = async {
430 tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
431 .expect("SIGTERM handler should install")
432 .recv()
433 .await;
434 };
435 #[cfg(not(unix))]
436 let terminate = std::future::pending::<()>();
437 tokio::select! {
438 () = ctrl_c => {}
439 _ = terminate => {}
440 }
441}