#if MYSQL_6_9 || MYSQL_6_10
/* 2021.10.14 */
using Externals.MySql.Data.MySqlClient;
using System;
using System.Collections.Generic;
using System.Data;
using System.Net;
using System.Text;
using System.Transactions;
using static Apewer.Source.OrmHelper;
namespace Apewer.Source
{
///
public class MySql : IDbClient
{
#region 基础
private Timeout _timeout = null;
private string _connectionstring = null;
/// 获取或设置日志记录。
public Logger Logger { get; set; }
/// 超时设定。
public Timeout Timeout { get => _timeout; }
/// 创建实例。
public MySql(string connnectionString, Timeout timeout = default)
{
_connectionstring = connnectionString;
_timeout = timeout ?? Timeout.Default;
}
/// 获取当前的 MySqlConnection 对象。
public IDbConnection Connection { get => _connection; }
/// 构建连接字符串以创建实例。
public MySql(string address, string store, string user, string pass, Timeout timeout = null)
{
_timeout = timeout ?? Timeout.Default;
var a = TextUtility.AntiInject(address);
var s = TextUtility.AntiInject(store);
var u = TextUtility.AntiInject(user);
var p = TextUtility.AntiInject(pass);
var cs = $"server={a}; database={s}; uid={u}; pwd={p}; ";
_connectionstring = cs;
_storename = new Class(s);
}
#endregion
#region 日志。
private void LogError(string action, Exception ex, string addtion)
{
var logger = Logger;
if (logger != null) logger.Error(this, "MySQL", action, ex.GetType().FullName, ex.Message, addtion);
}
#endregion
#region Connection
private MySqlConnection _connection = null;
///
public bool Online { get => _connection == null ? false : (_connection.State == ConnectionState.Open); }
/// 连接字符串。
public string ConnectionString { get => _connectionstring; }
///
public bool Connect()
{
if (_connection == null)
{
_connection = new MySqlConnection();
_connection.ConnectionString = _connectionstring;
}
else
{
if (_connection.State == ConnectionState.Open) return true;
}
// try
{
_connection.Open();
switch (_connection.State)
{
case ConnectionState.Open: return true;
default: return false;
}
}
// catch (Exception ex)
// {
// LogError("Connection", ex, _connection.ConnectionString);
// Close();
// return false;
// }
}
///
public void Close()
{
if (_connection != null)
{
if (_transaction != null)
{
if (_autocommit) Commit();
else Rollback();
}
_connection.Close();
_connection.Dispose();
_connection = null;
}
}
///
public void Dispose() { Close(); }
#endregion
#region Transaction
private IDbTransaction _transaction = null;
private bool _autocommit = false;
/// 启动事务。
public string Begin(bool commit = true) => Begin(commit, null);
/// 启动事务。
public string Begin(bool commit, Class isolation)
{
if (!Connect()) return "未连接。";
if (_transaction != null) return "存在已启动的事务,无法再次启动。";
try
{
_transaction = isolation ? _connection.BeginTransaction(isolation.Value) : _connection.BeginTransaction();
_autocommit = commit;
return null;
}
catch (Exception ex)
{
Logger.Error(nameof(MySql), "Commit", ex.Message());
return ex.Message();
}
}
/// 提交事务。
public string Commit()
{
if (_transaction == null) return "事务不存在。";
try
{
_transaction.Commit();
RuntimeUtility.Dispose(_transaction);
_transaction = null;
return null;
}
catch (Exception ex)
{
RuntimeUtility.Dispose(_transaction);
_transaction = null;
Logger.Error(nameof(MySql), "Commit", ex.Message());
return ex.Message();
}
}
/// 从挂起状态回滚事务。
public string Rollback()
{
if (_transaction == null) return "事务不存在。";
try
{
_transaction.Rollback();
RuntimeUtility.Dispose(_transaction);
_transaction = null;
return null;
}
catch (Exception ex)
{
RuntimeUtility.Dispose(_transaction);
_transaction = null;
Logger.Error(nameof(MySql), "Rollback", ex.Message());
return ex.Message();
}
}
#endregion
#region SQL
///
public IQuery Query(string sql, IEnumerable parameters)
{
if (sql.IsBlank()) return Example.InvalidQueryStatement;
var connected = Connect();
if (!connected) return Example.InvalidQueryConnection;
try
{
using (var command = new MySqlCommand())
{
command.Connection = _connection;
command.CommandTimeout = _timeout.Query;
command.CommandText = sql;
if (parameters != null)
{
foreach (var p in parameters)
{
if (p != null) command.Parameters.Add(p);
}
}
using (var ds = new DataSet())
{
using (var da = new MySqlDataAdapter(sql, _connection))
{
const string name = "result";
da.Fill(ds, name);
var table = ds.Tables[name];
return new Query(table);
}
}
}
}
catch (Exception exception)
{
Logger.Error(nameof(MySql), "Query", exception, sql);
return new Query(exception);
}
}
///
public IExecute Execute(string sql, IEnumerable parameters)
{
if (sql.IsBlank()) return Example.InvalidExecuteStatement;
var connected = Connect();
if (!connected) return Example.InvalidExecuteConnection;
var inTransaction = _transaction != null;
if (!inTransaction) Begin();
try
{
using (var command = new MySqlCommand())
{
command.Connection = _connection;
command.Transaction = (MySqlTransaction)_transaction;
command.CommandTimeout = _timeout.Execute;
command.CommandText = sql;
if (parameters != null)
{
foreach (var parameter in parameters)
{
if (parameter == null) continue;
command.Parameters.Add(parameter);
}
}
var rows = command.ExecuteNonQuery();
if (!inTransaction) Commit(); // todo 此处应该检查事务提交产生的错误。
return new Execute(true, rows);
}
}
catch (Exception exception)
{
Logger.Error(nameof(MySql), "Execute", exception, sql);
if (!inTransaction) Rollback();
return new Execute(exception);
}
}
///
public IQuery Query(string sql) => Query(sql, null);
///
public IExecute Execute(string sql, IEnumerable parameters)
{
var dps = null as List;
if (parameters != null)
{
var count = parameters.Count();
dps = new List(count);
foreach (var p in parameters)
{
var dp = CreateDataParameter(p);
dps.Add(dp);
}
}
return Execute(sql, dps);
}
///
public IExecute Execute(string sql) => Execute(sql, null as IEnumerable);
#endregion
#region ORM
private Class _storename = null;
private string StoreName()
{
if (_storename) return _storename.Value;
_storename = new Class(Internals.TextHelper.ParseConnectionString(_connectionstring).GetValue("database"));
return _storename.Value ?? "";
}
private string[] FirstColumn(string sql)
{
using (var query = Query(sql) as Query) return query.ReadColumn();
}
///
public string[] TableNames()
{
var sql = TextUtility.Merge("select table_name from information_schema.tables where table_schema='", StoreName(), "' and table_type='base table';");
return FirstColumn(sql);
}
///
public string[] ViewNames()
{
var sql = TextUtility.Merge("select table_name from information_schema.tables where table_schema='", StoreName(), "' and table_type='view';");
return FirstColumn(sql);
}
///
public string[] ColumnNames(string table)
{
var sql = TextUtility.Merge("select column_name from information_schema.columns where table_schema='", StoreName(), "' and table_name='", TextUtility.AntiInject(table), "';");
return FirstColumn(sql);
}
/// 获取用于创建表的语句。
private string GetCreateStetement(TableStructure structure)
{
// 检查现存表。
var exists = false;
var tables = TableNames();
if (tables.Length > 0)
{
var lower = structure.Name.ToLower();
foreach (var table in tables)
{
if (TextUtility.IsBlank(table)) continue;
if (table.ToLower() == lower)
{
exists = true;
break;
}
}
}
if (exists)
{
var columns = ColumnNames(structure.Name);
if (columns.Length > 0)
{
var lower = new List(columns.Length);
var added = 0;
foreach (var column in columns)
{
if (TextUtility.IsBlank(column)) continue;
lower.Add(column.ToLower());
added++;
}
lower.Capacity = added;
columns = lower.ToArray();
}
var sqlsb = new StringBuilder();
foreach (var column in structure.Columns)
{
// 检查 Independent 特性。
if (structure.Independent && column.Independent) continue;
// 去重。
var lower = column.Field.ToLower();
if (columns.Contains(lower)) continue;
var type = GetColumnDeclaration(column);
if (type.IsEmpty()) return TextUtility.Merge("类型 ", column.Type.ToString(), " 不受支持。");
// alter table `_record` add column `_index` bigint;
sqlsb.Append("alter table `", structure.Name, "` add column ", type, "; ");
}
var sql = sqlsb.ToString();
return sql;
}
else
{
// create table _record (`_index` bigint, `_key` varchar(255), `_text` longtext) engine=innodb default charset=utf8mb4
var columns = new List(structure.Columns.Length);
var columnsAdded = 0;
var primarykey = null as string;
foreach (var column in structure.Columns)
{
// 检查 Independent 特性。
if (structure.Independent && column.Independent) continue;
// 字段。
var type = GetColumnDeclaration(column);
if (type.IsEmpty()) return TextUtility.Merge("类型 ", column.Type.ToString(), " 不受支持。");
columns.Add(type);
columnsAdded++;
// 主键。
if (column.Property.Name == "Key") primarykey = column.Field;
}
columns.Capacity = columnsAdded;
var table = structure.Name;
var joined = string.Join(", ", columns);
// 设置主键。
string sql;
if (!structure.Independent && !string.IsNullOrEmpty(primarykey))
{
sql = TextUtility.Merge("create table `", table, "`(", joined, ", primary key (", primarykey, ") ) engine=innodb default charset=utf8mb4; ");
}
else
{
sql = TextUtility.Merge("create table `", table, "`(", joined, ") engine=innodb default charset=utf8mb4; ");
}
return sql;
}
}
///
private string Initialize(Type model, out string sql)
{
if (model == null)
{
sql = null;
return "指定的类型无效。";
}
var structure = TableStructure.Parse(model);
if (structure == null)
{
sql = null;
return "无法解析记录模型。";
}
// 连接数据库。
if (!Connect())
{
sql = null;
return "连接数据库失败。";
}
sql = GetCreateStetement(structure);
if (sql.NotEmpty())
{
var execute = Execute(sql);
if (!execute.Success) return execute.Message;
}
return null;
}
///
public string Initialize(Type model) => Initialize(model, out string sql);
///
public string Initialize() where T : class, new() => Initialize(typeof(T));
///
public string Initialize(Record model) => (model == null) ? "参数无效。" : Initialize(model.GetType());
/// 插入记录。返回错误信息。
public string Insert(object record)
{
if (record == null) return "参数无效。";
OrmHelper.FixProperties(record);
var structure = TableStructure.Parse(record.GetType());
if (structure == null) return "无法解析记录模型。";
var parameters = structure.CreateParameters(record, CreateDataParameter);
var sql = GenerateInsertStatement(structure.Name, parameters);
var execute = Execute(sql, parameters);
if (execute.Success) return TextUtility.Empty;
return execute.Message;
}
/// 更新记录,实体中的 Key 属性不被更新。返回错误信息。
/// 无法更新带有 Independent 特性的模型(缺少 Key 属性)。
public string Update(IRecord record)
{
if (record == null) return "参数无效。";
FixProperties(record);
SetUpdated(record);
var structure = TableStructure.Parse(record.GetType());
if (structure == null) return "无法解析记录模型。";
if (structure.Independent) return "无法更新带有 Independent 特性的模型。";
var parameters = structure.CreateParameters(record, CreateDataParameter, "_key");
var sql = GenerateUpdateStatement(structure, record.Key, parameters);
var execute = Execute(sql, parameters);
if (execute.Success) return TextUtility.Empty;
return execute.Message;
}
///
public Result