use std::panic::{AssertUnwindSafe, catch_unwind};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
#[cfg(feature = "metrics")]
use std::time::Instant;
use crate::fdb::tenant_name::FoundationDbTenantName;
use foundationdb::{Database, api::NetworkAutoStop, options::TransactionOption};
use crate::fdb::{
config::FdbConfig,
error::{FdbError, FdbHealthErrorReason, FdbResult, FdbSetupErrorReason},
tenant::{TenantHandle, open_provisioned},
};
#[cfg(feature = "metrics")]
use metrics::{counter, histogram};
static CLIENT_INITIALIZED: AtomicBool = AtomicBool::new(false);
#[cfg(feature = "metrics")]
const METRIC_HEALTH_CHECKS_TOTAL: &str = "reallyme_fdb_health_checks_total";
#[cfg(feature = "metrics")]
const METRIC_HEALTH_CHECK_LATENCY_SECONDS: &str = "reallyme_fdb_health_check_latency_seconds";
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct FdbHealthReport {
pub api_version: i32,
pub read_version: i64,
pub main_thread_busyness: Option<f64>,
}
struct FdbConnectorInner {
database: Database,
_network: NetworkAutoStop,
api_version: i32,
health_check_timeout: std::time::Duration,
}
#[derive(Clone)]
pub struct FoundationDbConnector {
inner: Arc<FdbConnectorInner>,
}
pub type FdbContext = FoundationDbConnector;
impl FoundationDbConnector {
pub fn connect(config: &FdbConfig) -> FdbResult<Self> {
CLIENT_INITIALIZED
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.map_err(|_| FdbError::Setup {
reason: FdbSetupErrorReason::ClientAlreadyInitialized,
})?;
let network_builder = catch_unwind(AssertUnwindSafe(|| {
foundationdb::api::FdbApiBuilder::default()
.set_runtime_version(config.fdb_api_version())
.build()
}))
.map_err(|_| FdbError::Setup {
reason: FdbSetupErrorReason::ClientInitializationPanicked,
})?
.map_err(|_| FdbError::Setup {
reason: FdbSetupErrorReason::ApiVersionUnsupported,
})?;
#[allow(unsafe_code)]
let network = catch_unwind(AssertUnwindSafe(|| unsafe { network_builder.boot() }))
.map_err(|_| FdbError::Setup {
reason: FdbSetupErrorReason::ClientInitializationPanicked,
})?
.map_err(|_| FdbError::Setup {
reason: FdbSetupErrorReason::NetworkBootFailed,
})?;
let database = match config.fdb_cluster_file() {
Some(cluster_file) => Database::from_path(cluster_file),
None => Database::default(),
}
.map_err(|_| FdbError::Setup {
reason: FdbSetupErrorReason::DatabaseOpenFailed,
})?;
Ok(Self {
inner: Arc::new(FdbConnectorInner {
_network: network,
database,
api_version: config.fdb_api_version(),
health_check_timeout: config.health_check_timeout(),
}),
})
}
pub async fn connect_and_check(config: &FdbConfig) -> FdbResult<Self> {
let connector = Self::connect(config)?;
connector.health_check().await?;
Ok(connector)
}
pub async fn health_check(&self) -> FdbResult<()> {
self.health_report().await.map(|_report| ())
}
pub async fn health_report(&self) -> FdbResult<FdbHealthReport> {
#[cfg(feature = "metrics")]
let started = Instant::now();
let probe = tokio::time::timeout(self.inner.health_check_timeout, self.probe()).await;
let result = match probe {
Ok(result) => result,
Err(_) => Err(FdbError::Health {
reason: FdbHealthErrorReason::DeadlineExceeded,
}),
};
#[cfg(feature = "metrics")]
record_health_metrics(&result, started.elapsed());
result
}
async fn probe(&self) -> FdbResult<FdbHealthReport> {
self.inner
.database
.perform_no_op()
.await
.map_err(|_| FdbError::Health {
reason: FdbHealthErrorReason::NetworkUnavailable,
})?;
let transaction = self
.inner
.database
.create_trx()
.map_err(|_| FdbError::Health {
reason: FdbHealthErrorReason::NetworkUnavailable,
})?;
transaction
.set_option(TransactionOption::AccessSystemKeys)
.map_err(|_| FdbError::Health {
reason: FdbHealthErrorReason::NetworkUnavailable,
})?;
let read_version = transaction
.get_read_version()
.await
.map_err(|_| FdbError::Health {
reason: FdbHealthErrorReason::ClusterUnavailable,
})?;
let main_thread_busyness = if self.inner.api_version >= 710 {
self.inner.database.get_main_thread_busyness().await.ok()
} else {
None
};
Ok(FdbHealthReport {
api_version: self.inner.api_version,
read_version,
main_thread_busyness,
})
}
pub async fn open_tenant(
&self,
tenant: FoundationDbTenantName,
) -> FdbResult<Arc<TenantHandle>> {
open_provisioned(&self.inner.database, tenant, self.clone()).await
}
#[allow(dead_code)]
pub(crate) fn database(&self) -> &Database {
&self.inner.database
}
#[cfg(feature = "tenant-admin")]
pub fn database_for_admin(&self) -> &Database {
&self.inner.database
}
}
impl std::fmt::Debug for FoundationDbConnector {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("FoundationDbConnector")
.field("api_version", &self.inner.api_version)
.field("health_check_timeout", &self.inner.health_check_timeout)
.field("database", &"<redacted-foundationdb-handle>")
.finish()
}
}
#[cfg(feature = "metrics")]
fn record_health_metrics(result: &FdbResult<FdbHealthReport>, duration: std::time::Duration) {
let result_label = if result.is_ok() { "ready" } else { "failed" };
counter!(METRIC_HEALTH_CHECKS_TOTAL, "result" => result_label).increment(1);
histogram!(METRIC_HEALTH_CHECK_LATENCY_SECONDS).record(duration.as_secs_f64());
}