mod settings;
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::EnvFilter;
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;
}
};
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 filter = match &settings.log {
Some(filter) => EnvFilter::new(filter.clone()),
None => EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")),
};
tracing_subscriber::fmt().with_env_filter(filter).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())?;
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(())
}