use std::any::Any;
use std::collections::HashMap;
use std::sync::Arc;
use camel_api::datasource::{
CheckFuture, CreatePoolFuture, DatasourceCatalog, DatasourceConfig, DatasourceHandle,
PoolFactory,
};
use camel_api::error::CamelError;
use camel_api::lifecycle::HealthStatus;
use camel_core::datasource::RuntimeDatasourceCatalog;
use crate::sql_action::{SqlAction, execute_sql_prepare};
const DB_URL: &str = "sqlite::memory:?cache=shared";
struct StubPoolFactory;
impl PoolFactory for StubPoolFactory {
fn create<'a>(&'a self, config: &'a DatasourceConfig) -> CreatePoolFuture<'a> {
Box::pin(async move {
sqlx::any::install_default_drivers();
let pool = sqlx::any::AnyPoolOptions::new()
.max_connections(1)
.connect(&config.db_url)
.await
.map_err(|e| CamelError::ProcessorError(e.to_string()))?;
Ok(Arc::new(pool) as Arc<dyn Any + Send + Sync>)
})
}
fn check<'a>(&'a self, _handle: &'a DatasourceHandle) -> CheckFuture<'a> {
Box::pin(async { HealthStatus::Healthy })
}
fn supported_schemes(&self) -> &[&str] {
&["sqlite"]
}
fn name(&self) -> &'static str {
"stub"
}
}
pub(crate) fn sqlite_catalog(name: &str) -> Arc<dyn DatasourceCatalog> {
let mut configs = HashMap::new();
configs.insert(
name.to_string(),
DatasourceConfig {
db_url: DB_URL.to_string(),
provider: None,
max_connections: None,
min_connections: None,
idle_timeout_secs: None,
max_lifetime_secs: None,
ssl_mode: None,
ssl_root_cert: None,
ssl_cert: None,
ssl_key: None,
extra: HashMap::new(),
},
);
let catalog = RuntimeDatasourceCatalog::new(configs);
assert!(
catalog
.register_factory("sqlite", Arc::new(StubPoolFactory))
.is_ok(),
"stub factory registration failed"
);
Arc::new(catalog)
}
pub(crate) async fn seed(catalog: &Arc<dyn DatasourceCatalog>, datasource: &str, stmts: &[&str]) {
let action = SqlAction {
datasource: datasource.to_string(),
prepare: stmts.iter().map(|stmt| stmt.to_string()).collect(),
};
if let Err(err) = execute_sql_prepare(catalog, &action).await {
panic!("seed failed: {err}");
}
}