use std::sync::Arc;
use fraiseql_core::{
cache::{CacheConfig, CachedDatabaseAdapter, QueryResultCache},
db::{
postgres::{
PoolPrewarmConfig, PostgresAdapter, PostgresTlsConfig, ReadReplicaPolicy, SearchPath,
},
traits::{DatabaseAdapter, Writer},
},
runtime::Executor,
schema::{CompiledSchema, TenancyMode},
};
use fraiseql_error::{FraiseQLError, Result};
use serde::Deserialize;
use tracing::info;
use super::schema_isolation;
#[derive(Debug, Clone, Deserialize)]
pub struct TenantPoolConfig {
pub connection_string: String,
#[serde(default = "default_max_connections")]
pub max_connections: u32,
#[serde(default = "default_connect_timeout")]
pub connect_timeout_secs: u64,
#[serde(default = "default_idle_timeout")]
pub idle_timeout_secs: u64,
#[serde(skip)]
pub search_path: Option<SearchPath>,
#[serde(skip)]
pub tls: PostgresTlsConfig,
#[serde(skip)]
pub vector_scan: fraiseql_core::db::postgres::VectorScanConfig,
#[serde(default)]
pub read_replica_urls: Vec<String>,
#[serde(skip)]
pub read_replica_policy: ReadReplicaPolicy,
}
const fn default_max_connections() -> u32 {
10
}
const fn default_connect_timeout() -> u64 {
5
}
const fn default_idle_timeout() -> u64 {
300
}
#[async_trait::async_trait]
pub trait FromPoolConfig: DatabaseAdapter + Sized {
async fn from_pool_config(config: &TenantPoolConfig) -> Result<Self>;
}
#[async_trait::async_trait]
impl FromPoolConfig for PostgresAdapter {
async fn from_pool_config(config: &TenantPoolConfig) -> Result<Self> {
Self::with_pool_config(
&config.connection_string,
PoolPrewarmConfig {
min_size: 0,
#[allow(clippy::cast_possible_truncation)]
max_size: config.max_connections as usize,
timeout_secs: Some(config.connect_timeout_secs),
search_path: config.search_path.clone(),
tls: config.tls.clone(),
read_replicas: config
.read_replica_policy
.with_urls(config.read_replica_urls.clone()),
max_streaming_reads: None,
vector_scan: config.vector_scan,
},
)
.await
}
}
#[async_trait::async_trait]
impl<A: FromPoolConfig> FromPoolConfig for CachedDatabaseAdapter<A> {
async fn from_pool_config(config: &TenantPoolConfig) -> Result<Self> {
let inner = A::from_pool_config(config).await?;
let cache = QueryResultCache::new(CacheConfig::default());
Ok(Self::new(inner, cache, String::new()))
}
}
#[doc(hidden)] pub async fn create_tenant_executor<A: FromPoolConfig + Writer>(
tenant_key: &str,
schema_json: &str,
pool_config: &TenantPoolConfig,
runtime_config: &fraiseql_core::runtime::RuntimeConfig,
) -> Result<Arc<Executor>> {
create_tenant_executor_with_adapter::<A>(tenant_key, schema_json, pool_config, runtime_config)
.await
.map(|(executor, _adapter)| executor)
}
#[doc(hidden)] pub async fn create_tenant_executor_with_adapter<A: FromPoolConfig + Writer>(
tenant_key: &str,
schema_json: &str,
pool_config: &TenantPoolConfig,
runtime_config: &fraiseql_core::runtime::RuntimeConfig,
) -> Result<(Arc<Executor>, Arc<A>)> {
let schema =
CompiledSchema::from_json(schema_json, false).map_err(|e| FraiseQLError::Parse {
message: format!("Invalid compiled schema JSON: {e}"),
location: String::new(),
})?;
schema
.validate_producer_version()
.map_err(|msg| FraiseQLError::validation(format!("Incompatible compiled schema: {msg}")))?;
let tenancy_mode = schema.tenancy_mode();
let pool_config = TenantPoolConfig {
search_path: if tenancy_mode == TenancyMode::Schema {
Some(schema_isolation::tenant_search_path(tenant_key)?)
} else {
None
},
..pool_config.clone()
};
let adapter = A::from_pool_config(&pool_config).await?;
if tenancy_mode == TenancyMode::Schema {
info!(tenant_key, "provisioning schema for tenant (schema isolation mode)");
schema_isolation::provision_tenant_schema(tenant_key, &adapter).await?;
schema_isolation::verify_search_path(tenant_key, &adapter).await?;
}
fraiseql_core::schema::refuse_standby_unreadable_sources(&adapter, &schema).await?;
let config = runtime_config
.clone()
.with_compiled_schema(&schema)
.map_err(|msg| FraiseQLError::validation(format!("Incompatible compiled schema: {msg}")))?;
let adapter = Arc::new(adapter);
Ok((Arc::new(Executor::with_config(schema, Arc::clone(&adapter), config)), adapter))
}
pub async fn destroy_tenant_schema(
tenant_key: &str,
executor: &fraiseql_core::runtime::Executor,
) -> Result<()> {
schema_isolation::drop_tenant_schema(tenant_key, executor).await
}