burncloud-database-client 0.1.1

High-level database client with migrations, pooling, and AI model management for BurnCloud
Documentation
// 数据库迁移管理系统

use burncloud_database_core::error::{DatabaseResult, DatabaseError};
use burncloud_database_core::{QueryExecutor, QueryContext, MigrationInfo};
use async_trait::async_trait;
use std::collections::HashMap;
use chrono::{DateTime, Utc};

/// 数据库迁移管理器
pub struct MigrationManager {
    query_executor: Box<dyn QueryExecutor>,
    migrations: Vec<Migration>,
}

/// 迁移定义
pub struct Migration {
    pub version: String,
    pub name: String,
    pub up_sql: String,
    pub down_sql: String,
    pub checksum: String,
}

impl MigrationManager {
    pub fn new(query_executor: Box<dyn QueryExecutor>) -> Self {
        let mut manager = Self {
            query_executor,
            migrations: Vec::new(),
        };

        // 注册所有迁移
        manager.register_migrations();
        manager
    }

    /// 注册所有数据库迁移
    fn register_migrations(&mut self) {
        // 1. 创建基础表结构
        self.add_migration(Migration {
            version: "001".to_string(),
            name: "create_base_tables".to_string(),
            up_sql: include_str!("../migrations/001_create_base_tables.sql").to_string(),
            down_sql: include_str!("../migrations/001_create_base_tables_down.sql").to_string(),
            checksum: self.calculate_checksum("001"),
        });

        // 2. AI模型管理表
        self.add_migration(Migration {
            version: "002".to_string(),
            name: "create_ai_model_tables".to_string(),
            up_sql: include_str!("../migrations/002_create_ai_model_tables.sql").to_string(),
            down_sql: include_str!("../migrations/002_create_ai_model_tables_down.sql").to_string(),
            checksum: self.calculate_checksum("002"),
        });

        // 3. 监控和指标表
        self.add_migration(Migration {
            version: "003".to_string(),
            name: "create_monitoring_tables".to_string(),
            up_sql: include_str!("../migrations/003_create_monitoring_tables.sql").to_string(),
            down_sql: include_str!("../migrations/003_create_monitoring_tables_down.sql").to_string(),
            checksum: self.calculate_checksum("003"),
        });

        // 4. 用户设置和安全表
        self.add_migration(Migration {
            version: "004".to_string(),
            name: "create_user_security_tables".to_string(),
            up_sql: include_str!("../migrations/004_create_user_security_tables.sql").to_string(),
            down_sql: include_str!("../migrations/004_create_user_security_tables_down.sql").to_string(),
            checksum: self.calculate_checksum("004"),
        });

        // 5. 索引和性能优化
        self.add_migration(Migration {
            version: "005".to_string(),
            name: "create_indexes".to_string(),
            up_sql: include_str!("../migrations/005_create_indexes.sql").to_string(),
            down_sql: include_str!("../migrations/005_create_indexes_down.sql").to_string(),
            checksum: self.calculate_checksum("005"),
        });
    }

    fn add_migration(&mut self, migration: Migration) {
        self.migrations.push(migration);
    }

    fn calculate_checksum(&self, version: &str) -> String {
        // 简单的校验和计算,实际应用中应使用更复杂的算法
        format!("checksum_{}", version)
    }

    /// 运行所有待执行的迁移
    pub async fn run_migrations(&self, context: &QueryContext) -> DatabaseResult<()> {
        // 确保迁移表存在
        self.create_migration_table(context).await?;

        // 获取已执行的迁移
        let applied_migrations = self.get_applied_migrations(context).await?;

        for migration in &self.migrations {
            if !applied_migrations.iter().any(|m| m.version == migration.version) {
                println!("Running migration: {} - {}", migration.version, migration.name);
                self.apply_migration(migration, context).await?;
                self.record_migration(migration, context).await?;
                println!("✅ Migration {} completed", migration.version);
            }
        }

        Ok(())
    }

    /// 回滚指定版本的迁移
    pub async fn rollback_migration(&self, version: &str, context: &QueryContext) -> DatabaseResult<()> {
        if let Some(migration) = self.migrations.iter().find(|m| m.version == version) {
            println!("Rolling back migration: {} - {}", migration.version, migration.name);

            // 执行回滚SQL
            self.execute_sql(&migration.down_sql, context).await?;

            // 从迁移记录中删除
            self.remove_migration_record(version, context).await?;

            println!("✅ Migration {} rolled back", version);
            Ok(())
        } else {
            Err(DatabaseError::ConfigurationError(format!("Migration {} not found", version)))
        }
    }

    /// 获取迁移状态
    pub async fn get_migration_status(&self, context: &QueryContext) -> DatabaseResult<Vec<MigrationInfo>> {
        self.get_applied_migrations(context).await
    }

    /// 创建迁移记录表
    async fn create_migration_table(&self, context: &QueryContext) -> DatabaseResult<()> {
        let sql = "
            CREATE TABLE IF NOT EXISTS schema_migrations (
                version VARCHAR(255) PRIMARY KEY,
                name VARCHAR(255) NOT NULL,
                applied_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(),
                checksum VARCHAR(255) NOT NULL
            )
        ";

        self.execute_sql(sql, context).await
    }

    /// 应用单个迁移
    async fn apply_migration(&self, migration: &Migration, context: &QueryContext) -> DatabaseResult<()> {
        self.execute_sql(&migration.up_sql, context).await
    }

    /// 记录已应用的迁移
    async fn record_migration(&self, migration: &Migration, context: &QueryContext) -> DatabaseResult<()> {
        let sql = "
            INSERT INTO schema_migrations (version, name, checksum)
            VALUES ($1, $2, $3)
        ";

        let version_param = burncloud_database_impl::StringParam(migration.version.clone());
        let name_param = burncloud_database_impl::StringParam(migration.name.clone());
        let checksum_param = burncloud_database_impl::StringParam(migration.checksum.clone());
        let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&version_param, &name_param, &checksum_param];

        self.query_executor.execute_query(sql, &params, context).await?;
        Ok(())
    }

    /// 删除迁移记录
    async fn remove_migration_record(&self, version: &str, context: &QueryContext) -> DatabaseResult<()> {
        let sql = "DELETE FROM schema_migrations WHERE version = $1";
        let version_param = burncloud_database_impl::StringParam(version.to_string());
        let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![&version_param];

        self.query_executor.execute_query(sql, &params, context).await?;
        Ok(())
    }

    /// 获取已应用的迁移
    async fn get_applied_migrations(&self, context: &QueryContext) -> DatabaseResult<Vec<MigrationInfo>> {
        let sql = "SELECT version, name, applied_at, checksum FROM schema_migrations ORDER BY version";
        let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![];

        let result = self.query_executor.execute_query(sql, &params, context).await?;

        let mut migrations = Vec::new();
        for row in result.rows {
            let version = row.get("version")
                .and_then(|v| v.as_str())
                .ok_or_else(|| DatabaseError::SerializationError("Missing version".to_string()))?
                .to_string();

            let name = row.get("name")
                .and_then(|v| v.as_str())
                .ok_or_else(|| DatabaseError::SerializationError("Missing name".to_string()))?
                .to_string();

            let applied_at_str = row.get("applied_at")
                .and_then(|v| v.as_str())
                .ok_or_else(|| DatabaseError::SerializationError("Missing applied_at".to_string()))?;

            let applied_at = DateTime::parse_from_rfc3339(applied_at_str)
                .map_err(|e| DatabaseError::SerializationError(e.to_string()))?
                .with_timezone(&Utc);

            let checksum = row.get("checksum")
                .and_then(|v| v.as_str())
                .ok_or_else(|| DatabaseError::SerializationError("Missing checksum".to_string()))?
                .to_string();

            migrations.push(MigrationInfo {
                version,
                name,
                applied_at,
                checksum,
            });
        }

        Ok(migrations)
    }

    /// 执行SQL语句
    async fn execute_sql(&self, sql: &str, context: &QueryContext) -> DatabaseResult<()> {
        let params: Vec<&dyn burncloud_database_core::QueryParam> = vec![];
        self.query_executor.execute_query(sql, &params, context).await?;
        Ok(())
    }
}

#[async_trait]
impl burncloud_database_core::MigrationManager for MigrationManager {
    async fn run_migrations(&self) -> DatabaseResult<()> {
        let context = QueryContext::default();
        self.run_migrations(&context).await
    }

    async fn rollback_migration(&self, version: &str) -> DatabaseResult<()> {
        let context = QueryContext::default();
        self.rollback_migration(version, &context).await
    }

    async fn get_migration_status(&self) -> DatabaseResult<Vec<MigrationInfo>> {
        let context = QueryContext::default();
        self.get_migration_status(&context).await
    }
}