use crate::{BaseRepository, QueryExecutor, QueryContext, QueryOptions};
use burncloud_database_core::error::DatabaseResult;
use crate::models::*;
use burncloud_database_core::error::DatabaseError;
use async_trait::async_trait;
use std::sync::Arc;
use chrono::{DateTime, Utc};
use uuid::Uuid;
pub struct AiModelRepository {
pub base: BaseRepository<AiModel>,
}
impl AiModelRepository {
pub fn new(query_executor: Box<dyn QueryExecutor>) -> Self {
Self {
base: BaseRepository::new(query_executor, "ai_models".to_string()),
}
}
pub async fn find_by_status(&self, status: ModelStatus, context: &QueryContext) -> DatabaseResult<Vec<AiModel>> {
let query = "SELECT * FROM ai_models WHERE status = $1";
let status_str = serde_json::to_string(&status)
.map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
let status_param = burncloud_database_impl::StringParam(status_str);
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&status_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
let mut models = Vec::new();
for row in result.rows {
let model: AiModel = serde_json::from_value(serde_json::Value::Object(
row.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
models.push(model);
}
Ok(models)
}
pub async fn find_by_type(&self, model_type: ModelType, context: &QueryContext) -> DatabaseResult<Vec<AiModel>> {
let query = "SELECT * FROM ai_models WHERE model_type = $1";
let type_str = serde_json::to_string(&model_type)
.map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
let type_param = burncloud_database_impl::StringParam(type_str);
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&type_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
let mut models = Vec::new();
for row in result.rows {
let model: AiModel = serde_json::from_value(serde_json::Value::Object(
row.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
models.push(model);
}
Ok(models)
}
pub async fn search(&self, query_text: &str, context: &QueryContext) -> DatabaseResult<Vec<AiModel>> {
let query = "SELECT * FROM ai_models WHERE name ILIKE $1 OR description ILIKE $1";
let search_pattern = format!("%{}%", query_text);
let search_param = burncloud_database_impl::StringParam(search_pattern);
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&search_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
let mut models = Vec::new();
for row in result.rows {
let model: AiModel = serde_json::from_value(serde_json::Value::Object(
row.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
models.push(model);
}
Ok(models)
}
}
pub struct ModelDeploymentRepository {
pub base: BaseRepository<ModelDeployment>,
}
impl ModelDeploymentRepository {
pub fn new(query_executor: Box<dyn QueryExecutor>) -> Self {
Self {
base: BaseRepository::new(query_executor, "model_deployments".to_string()),
}
}
pub async fn find_by_model_id(&self, model_id: Uuid, context: &QueryContext) -> DatabaseResult<Vec<ModelDeployment>> {
let query = "SELECT * FROM model_deployments WHERE model_id = $1";
let model_id_param = burncloud_database_impl::StringParam(model_id.to_string());
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&model_id_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
let mut deployments = Vec::new();
for row in result.rows {
let deployment: ModelDeployment = serde_json::from_value(serde_json::Value::Object(
row.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
deployments.push(deployment);
}
Ok(deployments)
}
pub async fn find_by_status(&self, status: DeploymentStatus, context: &QueryContext) -> DatabaseResult<Vec<ModelDeployment>> {
let query = "SELECT * FROM model_deployments WHERE status = $1";
let status_str = serde_json::to_string(&status)
.map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
let status_param = burncloud_database_impl::StringParam(status_str);
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&status_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
let mut deployments = Vec::new();
for row in result.rows {
let deployment: ModelDeployment = serde_json::from_value(serde_json::Value::Object(
row.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
deployments.push(deployment);
}
Ok(deployments)
}
pub async fn find_by_port(&self, port: u16, context: &QueryContext) -> DatabaseResult<Option<ModelDeployment>> {
let query = "SELECT * FROM model_deployments WHERE port = $1 LIMIT 1";
let port_param = burncloud_database_impl::I64Param(port as i64);
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&port_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
if result.rows.is_empty() {
Ok(None)
} else {
let deployment: ModelDeployment = serde_json::from_value(serde_json::Value::Object(
result.rows[0].iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
Ok(Some(deployment))
}
}
pub async fn get_running_deployments(&self, context: &QueryContext) -> DatabaseResult<Vec<ModelDeployment>> {
self.find_by_status(DeploymentStatus::Running, context).await
}
}
pub struct SystemMetricsRepository {
pub base: BaseRepository<SystemMetrics>,
}
impl SystemMetricsRepository {
pub fn new(query_executor: Box<dyn QueryExecutor>) -> Self {
Self {
base: BaseRepository::new(query_executor, "system_metrics".to_string()),
}
}
pub async fn find_by_time_range(
&self,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
context: &QueryContext,
) -> DatabaseResult<Vec<SystemMetrics>> {
let query = "SELECT * FROM system_metrics WHERE timestamp >= $1 AND timestamp <= $2 ORDER BY timestamp DESC";
let start_param = burncloud_database_impl::StringParam(start_time.to_rfc3339());
let end_param = burncloud_database_impl::StringParam(end_time.to_rfc3339());
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&start_param, &end_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
let mut metrics = Vec::new();
for row in result.rows {
let metric: SystemMetrics = serde_json::from_value(serde_json::Value::Object(
row.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
metrics.push(metric);
}
Ok(metrics)
}
pub async fn get_latest(&self, context: &QueryContext) -> DatabaseResult<Option<SystemMetrics>> {
let query = "SELECT * FROM system_metrics ORDER BY timestamp DESC LIMIT 1";
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
if result.rows.is_empty() {
Ok(None)
} else {
let metric: SystemMetrics = serde_json::from_value(serde_json::Value::Object(
result.rows[0].iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
Ok(Some(metric))
}
}
pub async fn cleanup_old_metrics(&self, before_time: DateTime<Utc>, context: &QueryContext) -> DatabaseResult<u64> {
let query = "DELETE FROM system_metrics WHERE timestamp < $1";
let time_param = burncloud_database_impl::StringParam(before_time.to_rfc3339());
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&time_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
Ok(result.rows_affected)
}
}
pub struct ModelMetricsRepository {
pub base: BaseRepository<ModelMetrics>,
}
impl ModelMetricsRepository {
pub fn new(query_executor: Box<dyn QueryExecutor>) -> Self {
Self {
base: BaseRepository::new(query_executor, "model_metrics".to_string()),
}
}
pub async fn find_by_deployment_id(
&self,
deployment_id: Uuid,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
context: &QueryContext,
) -> DatabaseResult<Vec<ModelMetrics>> {
let query = "SELECT * FROM model_metrics WHERE deployment_id = $1 AND timestamp >= $2 AND timestamp <= $3 ORDER BY timestamp DESC";
let deployment_param = burncloud_database_impl::StringParam(deployment_id.to_string());
let start_param = burncloud_database_impl::StringParam(start_time.to_rfc3339());
let end_param = burncloud_database_impl::StringParam(end_time.to_rfc3339());
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&deployment_param, &start_param, &end_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
let mut metrics = Vec::new();
for row in result.rows {
let metric: ModelMetrics = serde_json::from_value(serde_json::Value::Object(
row.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
metrics.push(metric);
}
Ok(metrics)
}
pub async fn get_latest_for_deployment(&self, deployment_id: Uuid, context: &QueryContext) -> DatabaseResult<Option<ModelMetrics>> {
let query = "SELECT * FROM model_metrics WHERE deployment_id = $1 ORDER BY timestamp DESC LIMIT 1";
let deployment_param = burncloud_database_impl::StringParam(deployment_id.to_string());
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&deployment_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
if result.rows.is_empty() {
Ok(None)
} else {
let metric: ModelMetrics = serde_json::from_value(serde_json::Value::Object(
result.rows[0].iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
Ok(Some(metric))
}
}
}
pub struct RequestLogRepository {
pub base: BaseRepository<RequestLog>,
}
impl RequestLogRepository {
pub fn new(query_executor: Box<dyn QueryExecutor>) -> Self {
Self {
base: BaseRepository::new(query_executor, "request_logs".to_string()),
}
}
pub async fn find_by_deployment_id(
&self,
deployment_id: Uuid,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
context: &QueryContext,
) -> DatabaseResult<Vec<RequestLog>> {
let query = "SELECT * FROM request_logs WHERE deployment_id = $1 AND timestamp >= $2 AND timestamp <= $3 ORDER BY timestamp DESC";
let deployment_param = burncloud_database_impl::StringParam(deployment_id.to_string());
let start_param = burncloud_database_impl::StringParam(start_time.to_rfc3339());
let end_param = burncloud_database_impl::StringParam(end_time.to_rfc3339());
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&deployment_param, &start_param, &end_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
let mut logs = Vec::new();
for row in result.rows {
let log: RequestLog = serde_json::from_value(serde_json::Value::Object(
row.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
logs.push(log);
}
Ok(logs)
}
pub async fn find_errors(
&self,
start_time: DateTime<Utc>,
end_time: DateTime<Utc>,
context: &QueryContext,
) -> DatabaseResult<Vec<RequestLog>> {
let query = "SELECT * FROM request_logs WHERE status_code >= 400 AND timestamp >= $1 AND timestamp <= $2 ORDER BY timestamp DESC";
let start_param = burncloud_database_impl::StringParam(start_time.to_rfc3339());
let end_param = burncloud_database_impl::StringParam(end_time.to_rfc3339());
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&start_param, &end_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
let mut logs = Vec::new();
for row in result.rows {
let log: RequestLog = serde_json::from_value(serde_json::Value::Object(
row.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
logs.push(log);
}
Ok(logs)
}
}
pub struct UserSettingsRepository {
pub base: BaseRepository<UserSettings>,
}
impl UserSettingsRepository {
pub fn new(query_executor: Box<dyn QueryExecutor>) -> Self {
Self {
base: BaseRepository::new(query_executor, "user_settings".to_string()),
}
}
pub async fn find_by_user_id(&self, user_id: &str, context: &QueryContext) -> DatabaseResult<Option<UserSettings>> {
let query = "SELECT * FROM user_settings WHERE user_id = $1 LIMIT 1";
let user_param = burncloud_database_impl::StringParam(user_id.to_string());
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&user_param];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
if result.rows.is_empty() {
Ok(None)
} else {
let settings: UserSettings = serde_json::from_value(serde_json::Value::Object(
result.rows[0].iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
Ok(Some(settings))
}
}
}
pub struct SecurityConfigRepository {
pub base: BaseRepository<SecurityConfig>,
}
impl SecurityConfigRepository {
pub fn new(query_executor: Box<dyn QueryExecutor>) -> Self {
Self {
base: BaseRepository::new(query_executor, "security_configs".to_string()),
}
}
pub async fn get_current(&self, context: &QueryContext) -> DatabaseResult<Option<SecurityConfig>> {
let query = "SELECT * FROM security_configs ORDER BY created_at DESC LIMIT 1";
let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![];
let result = self.base.query_executor.execute_query(query, ¶ms, context).await?;
if result.rows.is_empty() {
Ok(None)
} else {
let config: SecurityConfig = serde_json::from_value(serde_json::Value::Object(
result.rows[0].iter().map(|(k, v)| (k.clone(), v.clone())).collect()
)).map_err(|e| DatabaseError::SerializationError(e.to_string()))?;
Ok(Some(config))
}
}
}