use std::net::SocketAddr;
use std::sync::Arc;
use crate::config::ConfigurationRegistrant;
use crate::controller::RouteController;
use crate::data::cache::RedisStorage;
use crate::doc::DocumentationRegistrant;
use crate::env::AppEnvironment;
use crate::gateway::GatewayConnect;
use crate::job::JobRegistry;
#[derive(Debug, Clone)]
pub struct StartupOptions {
pub env_file: Option<String>,
pub enable_http: bool,
pub enable_jobs: bool,
pub enable_consumers: bool,
pub service_protocol: String,
pub service_weight: i32,
pub auth_type: String,
pub worker_max_execute_time_minutes: u64,
pub event_loop_max_execute_time_minutes: u64,
pub blocked_thread_check_interval_millis: u64,
pub worker_pool_size: usize,
pub event_loop_pool_size: usize,
}
impl Default for StartupOptions {
fn default() -> Self {
Self::default_options()
}
}
impl StartupOptions {
pub fn default_options() -> Self {
Self {
env_file: None,
enable_http: true,
enable_jobs: true,
enable_consumers: false,
service_protocol: "http".to_string(),
service_weight: 1,
auth_type: "token".to_string(),
worker_max_execute_time_minutes: 2,
event_loop_max_execute_time_minutes: 1,
blocked_thread_check_interval_millis: 750,
worker_pool_size: 20,
event_loop_pool_size: 16,
}
}
pub fn with_env_file(mut self, path: impl Into<String>) -> Self {
self.env_file = Some(path.into());
self
}
pub fn with_http(mut self, enabled: bool) -> Self {
self.enable_http = enabled;
self
}
pub fn with_jobs(mut self, enabled: bool) -> Self {
self.enable_jobs = enabled;
self
}
pub fn with_consumers(mut self, enabled: bool) -> Self {
self.enable_consumers = enabled;
self
}
}
pub struct GenericStartup {
pub options: StartupOptions,
pub mount_paths: Vec<String>,
pub registrant: Option<Arc<ConfigurationRegistrant>>,
pub job_registry: JobRegistry,
pub gateway: Option<Arc<GatewayConnect>>,
pub redis: Option<RedisStorage>,
pub server_addr: Option<SocketAddr>,
pub socket_addr: Option<SocketAddr>,
pub consumer_names: Vec<String>,
docs_built: bool,
}
impl GenericStartup {
pub fn new(options: StartupOptions) -> Self {
Self {
options,
mount_paths: Vec::new(),
registrant: None,
job_registry: JobRegistry::new(),
gateway: None,
redis: None,
server_addr: None,
socket_addr: None,
consumer_names: Vec::new(),
docs_built: false,
}
}
pub fn with_mount_paths(mut self, paths: Vec<String>) -> Self {
self.mount_paths = paths;
self
}
pub async fn init(&mut self) -> anyhow::Result<()> {
AppEnvironment::with_env_file(self.options.env_file.as_deref())?;
Self::init_tracing();
Self::log_pool_options(&self.options);
let env = AppEnvironment::get();
let addr: SocketAddr = format!("0.0.0.0:{}", env.server_port).parse()?;
let socket_addr: SocketAddr = format!("0.0.0.0:{}", env.socket_port).parse()?;
self.registrant = Some(Arc::new(ConfigurationRegistrant::new(addr)));
self.server_addr = Some(addr);
self.socket_addr = Some(socket_addr);
if let Ok(mut reg) = DocumentationRegistrant::global().write() {
reg.bind_base_list(&self.mount_paths);
}
match RedisStorage::from_env() {
Ok(storage) => {
tracing::info!(
component = "cache",
backend = "redis",
"Redis storage initialized"
);
self.redis = Some(storage);
}
Err(e) => {
tracing::warn!(component = "cache", backend = "redis", error = %e, "Redis unavailable; continuing without cache");
}
}
Ok(())
}
fn init_tracing() {
use tracing_subscriber::EnvFilter;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
let (env_filter, filter_source) = match EnvFilter::try_from_default_env() {
Ok(filter) => (filter, "RUST_LOG"),
Err(_) => (EnvFilter::new("INFO"), "DEFAULT"),
};
let _ = tracing_subscriber::registry()
.with(env_filter)
.with(tracing_subscriber::fmt::layer().with_ansi(true))
.try_init();
tracing::info!(filter_source, "Tracing initialized");
}
fn log_pool_options(options: &StartupOptions) {
if options.worker_pool_size == 0 || options.event_loop_pool_size == 0 {
tracing::warn!(
worker_pool_size = options.worker_pool_size,
event_loop_pool_size = options.event_loop_pool_size,
"Invalid runtime pool-size hints; Tokio runtime sizing is owned by the host binary"
);
} else {
tracing::info!(
worker_max_execute_time_minutes = options.worker_max_execute_time_minutes,
event_loop_max_execute_time_minutes = options.event_loop_max_execute_time_minutes,
blocked_thread_check_interval_ms = options.blocked_thread_check_interval_millis,
worker_pool_size = options.worker_pool_size,
event_loop_pool_size = options.event_loop_pool_size,
"Configured runtime pool-size hints"
);
}
}
pub async fn bootstrap(
options: StartupOptions,
mount_paths: Vec<String>,
static_registrar: Option<Arc<dyn StaticRegistrar>>,
controller_registrar: Option<Arc<dyn ControllerRegistrar>>,
consumer_registrar: Option<Arc<dyn ConsumerRegistrar>>,
) -> anyhow::Result<Self> {
let mut startup = Self::new(options);
startup.mount_paths = mount_paths.clone();
startup.init().await?;
let env = AppEnvironment::get();
let cpu_count = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1);
if let Some(sr) = static_registrar.as_ref() {
if let Some(reg) = startup.registrant.as_ref() {
let handle = reg.router_handle();
let mut router = handle.write().await;
sr.register_static(&startup, &mut router);
}
}
if let Ok(mut reg) = DocumentationRegistrant::global().write() {
reg.bind_base_list(&mount_paths);
}
match GatewayConnect::set_up(
mount_paths.clone(),
startup.options.service_protocol.clone(),
startup.options.service_weight,
startup.options.auth_type.clone(),
)
.await
{
Ok(Some(gw)) => {
tracing::info!(
component = "gateway",
mount_path_count = mount_paths.len(),
"Registered service configuration"
);
startup.gateway = Some(gw);
}
Ok(None) => {
tracing::info!(
component = "gateway",
"No mount paths configured; skipped service registration"
);
}
Err(e) => {
tracing::error!(component = "gateway", error = %e, "Failed to register service configuration");
if !env.is_production() {
eprintln!("{:?}", e);
}
}
}
let deployable_count = env.server_count + env.socket_count + env.worker_count;
if deployable_count > cpu_count {
tracing::warn!(
deployable_count,
cpu_count,
"Configured deployables exceed available CPU cores"
);
}
if env.server_count > 0 && startup.options.enable_http {
if let Some(cr) = controller_registrar.as_ref() {
let controllers = cr.controllers(&startup);
for ctrl in controllers {
startup.mount_controller_boxed(ctrl).await?;
}
}
startup
.mount_controller_boxed(Box::new(crate::doc::controller::DocumentationController))
.await?;
let timer = std::time::Instant::now();
crate::doc::controller::DocumentationController::build_specs();
startup.docs_built = true;
tracing::info!(
component = "openapi",
duration_ms = timer.elapsed().as_millis() as u64,
"Built OpenAPI documentation"
);
}
if env.worker_count > 0 {
tracing::info!(
worker_count = env.worker_count,
worker_pool_size = startup.options.worker_pool_size,
"Jobs are available via run_jobs()"
);
}
if let Some(cons_reg) = consumer_registrar.as_ref() {
if startup.options.enable_consumers {
let consumers = cons_reg.consumers(&startup);
for c in &consumers {
tracing::info!(queue = %c.queue_name(), "Registered consumer");
startup.consumer_names.push(c.queue_name().to_string());
}
if consumers.is_empty() {
tracing::info!("Consumer registrar provided no consumers");
}
} else {
tracing::info!(enabled = false, "Consumer registration skipped");
}
}
if env.socket_count > 0 {
tracing::info!(
socket_count = env.socket_count,
socket_addr = ?startup.socket_addr,
"Socket server is available for binding"
);
}
Ok(startup)
}
pub async fn mount_controller<C: RouteController + 'static>(&self, c: C) -> anyhow::Result<()> {
if let Some(r) = &self.registrant {
r.mount_controller(c).await;
Ok(())
} else {
anyhow::bail!("not initialized — call init() or bootstrap() first")
}
}
pub async fn mount_controller_boxed(&self, c: Box<dyn RouteController>) -> anyhow::Result<()> {
if let Some(r) = &self.registrant {
let handle = r.router_handle();
let mut router = handle.write().await;
tracing::info!(
target: "routing",
handler = c.type_name(),
path = c.base_path(),
"Mounted controller '{}' at '{}'",
c.type_name(),
c.base_path()
);
c.register_routes(&mut router).await;
Ok(())
} else {
anyhow::bail!("not initialized")
}
}
pub async fn serve(&self) -> anyhow::Result<SocketAddr> {
if let Some(r) = self.registrant.clone() {
Ok(r.serve().await?)
} else {
anyhow::bail!("not initialized")
}
}
pub async fn shutdown(&mut self) -> anyhow::Result<()> {
if let Err(e) = self.job_registry.stop().await {
tracing::warn!(component = "jobs", error = %e, "Failed to stop jobs during shutdown");
}
if let Some(gw) = self.gateway.take() {
if let Err(e) = gw.deregister().await {
tracing::warn!(component = "gateway", error = %e, "Failed to deregister service during shutdown");
}
}
tracing::info!("Shutdown complete");
Ok(())
}
pub fn job_registry(&mut self) -> &mut JobRegistry {
&mut self.job_registry
}
pub fn add_job<J: crate::job::ServiceJob + 'static>(&mut self, job: J) {
self.job_registry.add_job(job);
}
pub async fn run_jobs(&mut self) -> anyhow::Result<()> {
if !self.options.enable_jobs {
tracing::info!(enabled = false, "Job startup skipped");
return Ok(());
}
if self.job_registry.job_count() == 0 {
tracing::info!("No jobs registered; nothing to start");
return Ok(());
}
tracing::info!(job_count = self.job_registry.job_count(), "Starting jobs");
self.job_registry.start().await?;
Ok(())
}
pub async fn start_jobs(&mut self) -> anyhow::Result<()> {
self.run_jobs().await
}
pub async fn stop_jobs(&mut self) -> anyhow::Result<()> {
self.job_registry.stop().await?;
Ok(())
}
pub fn registrant(&self) -> Option<Arc<ConfigurationRegistrant>> {
self.registrant.clone()
}
pub fn gateway_client(&self) -> Option<Arc<GatewayConnect>> {
self.gateway.clone()
}
pub fn redis_client(&self) -> Option<RedisStorage> {
self.redis.clone()
}
pub fn server_addr(&self) -> Option<SocketAddr> {
self.server_addr
}
pub fn socket_addr(&self) -> Option<SocketAddr> {
self.socket_addr
}
pub fn router_handle(&self) -> Option<Arc<tokio::sync::RwLock<crate::controller::Router>>> {
self.registrant.as_ref().map(|r| r.router_handle())
}
pub fn is_production(&self) -> bool {
AppEnvironment::try_get()
.map(|e| e.is_production())
.unwrap_or(false)
}
pub fn docs_built(&self) -> bool {
self.docs_built
}
pub fn server_urls(&self) -> Vec<String> {
let port = AppEnvironment::try_get()
.map(|e| e.server_port)
.unwrap_or(8080);
vec![
format!("http://localhost:{port}"),
format!("http://127.0.0.1:{port}"),
]
}
}
pub trait ControllerRegistrar: Send + Sync {
fn controllers(&self, startup: &GenericStartup) -> Vec<Box<dyn RouteController>>;
}
pub trait ConsumerRegistrar: Send + Sync {
fn consumers(
&self,
startup: &GenericStartup,
) -> Vec<Box<dyn queue_descriptor::QueueDescriptor>>;
}
pub trait StaticRegistrar: Send + Sync {
fn register_static(&self, startup: &GenericStartup, router: &mut crate::controller::Router);
}
pub mod queue_descriptor {
pub trait QueueDescriptor: Send + Sync {
fn queue_name(&self) -> &str;
}
}