mod healthcheck;
mod settings;
use std::env;
use std::error::Error;
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 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 {
return match healthcheck::probe(&settings.bind) {
Ok(()) => ExitCode::SUCCESS,
Err(err) => {
eprintln!("healthcheck failed: {err}");
ExitCode::FAILURE
}
};
}
init_tracing(&settings);
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("RUST_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,
};
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(())
}