using System;
using System.Collections.Generic;
using System.Data;
using System.Text;
using System.Threading.Tasks;
using MySqlConnector;
using TeamAAS.Database.Attributes;
using TeamAAS.Database.Interfaces;
using TeamAAS.Database.Models;
namespace TeamAAS.Database.Providers
{
///
/// MySQL 数据库提供者实现
///
[DatabaseProvider("mysql", DisplayName = "MySQL", RequiresServer = true,
Description = "MySQL / MariaDB 数据库,基于 MySqlConnector")]
public class MySqlDatabase : IDatabase
{
private MySqlConnection _connection;
private DatabaseConfig _config;
private readonly object _sync = new object();
public Guid Id { get; private set; }
public string Name { get; set; }
public string ProviderType => "mysql";
public bool IsConnected => _connection != null && _connection.State == ConnectionState.Open;
public string ConnectionString
{
get
{
if (_config == null) return string.Empty;
if (!string.IsNullOrWhiteSpace(_config.ConnectionString))
return MaskPassword(_config.ConnectionString);
var sb = new StringBuilder();
sb.Append($"Server={_config.Server};");
if (_config.Port > 0) sb.Append($"Port={_config.Port};");
sb.Append($"Database={_config.DatabaseName};");
sb.Append($"Uid={_config.UserId};");
sb.Append("Pwd=***;");
if (_config.ConnectionTimeout > 0)
sb.Append($"Connection Timeout={_config.ConnectionTimeout};");
return sb.ToString();
}
}
private string GetRawConnectionString()
{
if (!string.IsNullOrWhiteSpace(_config?.ConnectionString))
return _config.ConnectionString;
if (_config == null)
throw new InvalidOperationException("数据库未配置");
var sb = new MySqlConnectionStringBuilder
{
Server = _config.Server,
Port = (uint)(_config.Port > 0 ? _config.Port : 3306),
Database = _config.DatabaseName,
UserID = _config.UserId,
Password = _config.Password,
ConnectionTimeout = (uint)(_config.ConnectionTimeout > 0 ? _config.ConnectionTimeout : 30)
};
if (_config.ExtraParams != null)
{
foreach (var kv in _config.ExtraParams)
{
try { sb[kv.Key] = kv.Value; } catch { /* 忽略无效参数 */ }
}
}
return sb.ConnectionString;
}
public void Configure(DatabaseConfig config)
{
_config = config ?? throw new ArgumentNullException(nameof(config));
Id = config.Id == Guid.Empty ? Guid.NewGuid() : config.Id;
Name = config.Name ?? $"MySQL_{_config.Server}";
}
public bool Open()
{
lock (_sync)
{
try
{
if (_connection == null)
_connection = new MySqlConnection(GetRawConnectionString());
if (_connection.State == ConnectionState.Open)
return true;
_connection.Open();
return true;
}
catch
{
return false;
}
}
}
public Task OpenAsync()
{
return Task.Run(() => Open());
}
public void Close()
{
lock (_sync)
{
if (_connection != null)
{
try { _connection.Close(); } catch { }
}
}
}
public Task CloseAsync()
{
return Task.Run(() => Close());
}
public (bool Success, string Message) TestConnection()
{
try
{
using (var conn = new MySqlConnection(GetRawConnectionString()))
{
conn.Open();
return (true, "连接成功");
}
}
catch (Exception ex)
{
return (false, ex.Message);
}
}
public Task<(bool Success, string Message)> TestConnectionAsync()
{
return Task.Run(() => TestConnection());
}
public int ExecuteNonQuery(string sql, IDictionary parameters = null)
{
using (var cmd = CreateCommand(sql, parameters))
{
EnsureOpen();
return cmd.ExecuteNonQuery();
}
}
public Task ExecuteNonQueryAsync(string sql, IDictionary parameters = null)
{
return Task.Run(() => ExecuteNonQuery(sql, parameters));
}
public DataTable ExecuteQuery(string sql, IDictionary parameters = null)
{
using (var cmd = CreateCommand(sql, parameters))
{
EnsureOpen();
var dt = new DataTable();
using (var reader = cmd.ExecuteReader())
{
dt.Load(reader);
}
return dt;
}
}
public Task ExecuteQueryAsync(string sql, IDictionary parameters = null)
{
return Task.Run(() => ExecuteQuery(sql, parameters));
}
public object ExecuteScalar(string sql, IDictionary parameters = null)
{
using (var cmd = CreateCommand(sql, parameters))
{
EnsureOpen();
return cmd.ExecuteScalar();
}
}
public Task