use dataflow_rs::datalogic_rs;
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use datalogic_rs::Engine as DatalogicEngine;
use metrics_exporter_prometheus::PrometheusHandle;
use tokio::sync::Mutex;
use crate::channel::ChannelRegistry;
use crate::config::AppConfig;
use crate::connector::ConnectorRegistry;
use crate::connector::cache_backend::CachePool;
use crate::queue::TraceQueue;
use crate::server::rate_limit::RateLimitState;
use crate::storage::DbPool;
use crate::storage::repositories::Repositories;
pub struct Kafka {
pub producer: Option<Arc<crate::kafka::producer::KafkaProducer>>,
pub consumer_handle: Arc<Mutex<Option<crate::kafka::consumer::ConsumerHandle>>>,
pub ingest_status: Arc<crate::kafka::KafkaIngestStatus>,
}
pub struct Caches {
pub cache_pool: Arc<CachePool>,
pub sql_pool_cache: Arc<crate::connector::pool_cache::SqlPoolCache>,
pub mongo_pool_cache: Arc<crate::connector::mongo_pool::MongoPoolCache>,
}
pub struct AppStateInner {
pub engine: Arc<crate::engine::EngineHandle>,
pub repos: Repositories,
pub audit_queue: crate::queue::audit_queue::AuditQueue,
pub connector_registry: Arc<ConnectorRegistry>,
pub caches: Caches,
pub channel_registry: Arc<ChannelRegistry>,
pub trace_queue: TraceQueue,
#[doc(hidden)]
pub db_pool: DbPool,
pub config: Arc<AppConfig>,
pub start_time: chrono::DateTime<chrono::Utc>,
pub metrics_handle: PrometheusHandle,
pub http_client: reqwest::Client,
pub datalogic: Arc<DatalogicEngine>,
pub rate_limit_state: Option<Arc<RateLimitState>>,
pub ready: Arc<AtomicBool>,
pub kafka: Kafka,
pub trace_persistence_queue: crate::queue::TracePersistenceQueue,
pub cluster: Arc<crate::cluster::ClusterRuntime>,
pub admin_auth_failures: Arc<crate::server::admin_auth::FailedAuthTracker>,
pub trusted_proxies: Arc<Vec<ipnet::IpNet>>,
}
impl AppStateInner {
pub fn trusted_proxies(&self) -> &[ipnet::IpNet] {
&self.trusted_proxies
}
pub fn pool_stats(&self) -> (u32, usize) {
(self.db_pool.size(), self.db_pool.num_idle())
}
pub async fn ping_db(&self) -> Result<(), sqlx::Error> {
self.db_pool.ping().await
}
pub async fn backup_sqlite_into(&self, path: &str) -> Result<bool, sqlx::Error> {
let DbPool::Sqlite(pool) = &self.db_pool else {
return Ok(false);
};
sqlx::query(&format!("VACUUM INTO '{}'", path.replace('\'', "''")))
.execute(pool)
.await?;
Ok(true)
}
}
pub type AppState = Arc<AppStateInner>;