mod healthcheck;
mod settings;
#[cfg(feature = "jemalloc")]
#[global_allocator]
static GLOBAL: tikv_jemallocator::Jemalloc = tikv_jemallocator::Jemalloc;
use std::env;
use std::error::Error;
use std::path::Path;
use std::process::{self, ExitCode};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use tephra::log::set::SegmentSet;
use tephra::writer::WriteCoordinator;
use tephra_server::Server;
use tephra_server::auth::AuthConfig;
use tracing_subscriber::filter::{LevelFilter, Targets};
use tracing_subscriber::layer::SubscriberExt as _;
use tracing_subscriber::util::SubscriberInitExt as _;
use settings::{Args, Settings};
const EXIT_FORCED: i32 = 143;
fn main() -> ExitCode {
let args: Args = argh::from_env();
let settings = match settings::load(&args) {
Ok(settings) => settings,
Err(err) => {
eprintln!("configuration error: {err}");
return ExitCode::FAILURE;
}
};
if args.healthcheck {
let tls_cert = settings.tls.cert.as_deref().map(Path::new);
return match healthcheck::probe(&settings.bind, settings.first_auth_token(), tls_cert) {
Ok(()) => ExitCode::SUCCESS,
Err(err) => {
eprintln!("healthcheck failed: {err}");
ExitCode::FAILURE
}
};
}
init_tracing(&settings);
#[cfg(feature = "jemalloc")]
if let Err(err) = tikv_jemalloc_ctl::background_thread::write(true) {
tracing::warn!(%err, "failed to enable jemalloc background threads");
}
match run(settings) {
Ok(()) => ExitCode::SUCCESS,
Err(err) => {
tracing::error!(%err, "server exited with error");
ExitCode::FAILURE
}
}
}
fn init_tracing(settings: &Settings) {
let directives = settings
.log
.clone()
.or_else(|| env::var("TEPHRA_LOG").ok())
.unwrap_or_else(|| "info".to_string());
let targets = directives
.parse::<Targets>()
.unwrap_or_else(|_| Targets::new().with_default(LevelFilter::INFO));
tracing_subscriber::registry()
.with(tracing_subscriber::fmt::layer())
.with(targets)
.init();
}
fn run(settings: Settings) -> Result<(), Box<dyn Error>> {
let set = SegmentSet::open(&settings.data_dir, settings.segment_config())?;
let capacity = set.segment_capacity();
let mut writer_config = settings.writer_config();
if writer_config.max_batch_bytes > capacity {
tracing::warn!(
requested = writer_config.max_batch_bytes,
capacity,
"max_batch_bytes exceeds segment capacity; clamping to capacity"
);
writer_config.max_batch_bytes = capacity;
}
let (coordinator, handle) = WriteCoordinator::start(set, writer_config)?;
tracing::info!(
data_dir = %settings.data_dir,
segment_size = settings.segment.size,
"opened event store"
);
let server = Server::bind(&settings.bind, handle, settings.server_config())?
.with_data_dir(&settings.data_dir);
#[cfg(feature = "metrics")]
let server = match &settings.metrics.bind {
Some(addr) => server.with_metrics_addr(addr)?,
None => server,
};
#[cfg(feature = "tls")]
let server = match (&settings.tls.cert, &settings.tls.key) {
(Some(cert), Some(key)) => {
let tls = tephra_server::tls::build_server_config(Path::new(cert), Path::new(key))?;
tracing::info!("serving over tls");
server.with_tls(tls)
}
_ => server,
};
#[cfg(not(feature = "tls"))]
if settings.tls.cert.is_some() {
tracing::warn!("tls is configured but this binary was built without the tls feature");
}
let server = match settings.auth_tokens() {
Some(tokens) => {
tracing::info!(
tokens = tokens.len(),
"requiring bearer-token authentication"
);
server.with_auth(Arc::new(AuthConfig::new(tokens)))
}
None => server,
};
let shutdown = server.shutdown_handle();
let signalled = Arc::new(AtomicBool::new(false));
ctrlc::set_handler(move || {
if signalled.swap(true, Ordering::SeqCst) {
tracing::warn!("received second signal, exiting immediately");
process::exit(EXIT_FORCED);
}
tracing::info!("received shutdown signal, shutting down");
shutdown.shutdown();
})?;
let run_result = server.run();
coordinator.shutdown();
run_result?;
Ok(())
}