use burncloud_database_core::{
DatabaseConnection, QueryExecutor, TransactionManager, Transaction,
QueryContext, QueryOptions, QueryResult, QueryParam
};
use burncloud_database_core::error::{DatabaseResult, DatabaseError};
use async_trait::async_trait;
use std::sync::Arc;
use tokio::sync::RwLock;
use tracing::{info, error, debug};
pub struct DatabaseClient {
connection: Arc<RwLock<Box<dyn DatabaseConnection>>>,
query_executor: Arc<Box<dyn QueryExecutor>>,
}
impl DatabaseClient {
pub fn new(
connection: Box<dyn DatabaseConnection>,
query_executor: Box<dyn QueryExecutor>,
) -> Self {
Self {
connection: Arc::new(RwLock::new(connection)),
query_executor: Arc::new(query_executor),
}
}
pub async fn connect(&self) -> DatabaseResult<()> {
let mut conn = self.connection.write().await;
info!("Connecting to database...");
conn.connect().await?;
info!("Successfully connected to database");
Ok(())
}
pub async fn disconnect(&self) -> DatabaseResult<()> {
let mut conn = self.connection.write().await;
info!("Disconnecting from database...");
conn.disconnect().await?;
info!("Successfully disconnected from database");
Ok(())
}
pub async fn is_connected(&self) -> bool {
let conn = self.connection.read().await;
conn.is_connected().await
}
pub async fn ping(&self) -> DatabaseResult<()> {
let conn = self.connection.read().await;
debug!("Pinging database...");
conn.ping().await?;
debug!("Database ping successful");
Ok(())
}
pub async fn execute_query(
&self,
query: &str,
params: &[&dyn QueryParam],
context: &QueryContext,
) -> DatabaseResult<QueryResult> {
debug!("Executing query: {}", query);
let result = self.query_executor.execute_query(query, params, context).await;
match &result {
Ok(res) => debug!("Query executed successfully, {} rows returned", res.rows.len()),
Err(e) => error!("Query execution failed: {}", e),
}
result
}
pub async fn execute_query_with_options(
&self,
query: &str,
params: &[&dyn QueryParam],
options: &QueryOptions,
context: &QueryContext,
) -> DatabaseResult<QueryResult> {
debug!("Executing query with options: {}", query);
let result = self.query_executor.execute_query_with_options(query, params, options, context).await;
match &result {
Ok(res) => debug!("Query with options executed successfully, {} rows returned", res.rows.len()),
Err(e) => error!("Query with options execution failed: {}", e),
}
result
}
}
#[async_trait]
impl DatabaseConnection for DatabaseClient {
async fn connect(&mut self) -> DatabaseResult<()> {
self.connect().await
}
async fn disconnect(&mut self) -> DatabaseResult<()> {
self.disconnect().await
}
async fn is_connected(&self) -> bool {
self.is_connected().await
}
async fn ping(&self) -> DatabaseResult<()> {
self.ping().await
}
}
#[async_trait]
impl QueryExecutor for DatabaseClient {
async fn execute_query(
&self,
query: &str,
params: &[&dyn QueryParam],
context: &QueryContext,
) -> DatabaseResult<QueryResult> {
self.execute_query(query, params, context).await
}
async fn execute_query_with_options(
&self,
query: &str,
params: &[&dyn QueryParam],
options: &QueryOptions,
context: &QueryContext,
) -> DatabaseResult<QueryResult> {
self.execute_query_with_options(query, params, options, context).await
}
}