ProductDatabase.cs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525
  1. using System;
  2. using System.Collections.Generic;
  3. using System.Data;
  4. using System.IO;
  5. using TeamAAS.Database.Interfaces;
  6. using TeamAAS.Database.Models;
  7. using TeamAAS.Database.Providers;
  8. namespace TeamAAS.Database.Services
  9. {
  10. /// <summary>
  11. /// 产品数据库(生产记录 / 运行日志 / 报警历史 / 设备运行数据)。
  12. /// 跟着产品配方走,每产品一个独立数据库文件;也支持配置为 MySQL 远程库。
  13. /// </summary>
  14. public class ProductDatabase : IDisposable
  15. {
  16. private readonly IDatabase _db;
  17. private readonly DatabaseConfig _config;
  18. /// <summary>
  19. /// 数据库实例
  20. /// </summary>
  21. public IDatabase Database => _db;
  22. /// <summary>当前数据库配置(供配置页展示/编辑;改动后由 ProductDatabaseManager 重建实例生效)。</summary>
  23. public DatabaseConfig Config => _config;
  24. /// <summary>
  25. /// 是否连接成功
  26. /// </summary>
  27. public bool IsConnected => _db?.IsConnected == true;
  28. /// <summary>
  29. /// 当前产品名
  30. /// </summary>
  31. public string ProductName { get; private set; }
  32. /// <summary>
  33. /// 用配置创建产品数据库(支持 SQLite / MySQL)
  34. /// </summary>
  35. public ProductDatabase(DatabaseConfig config)
  36. {
  37. _config = config ?? throw new ArgumentNullException(nameof(config));
  38. ProductName = config.Name;
  39. // 根据 ProviderType 选择实现
  40. switch ((config.ProviderType ?? "sqlite").ToLowerInvariant())
  41. {
  42. case "mysql":
  43. _db = new MySqlDatabase();
  44. break;
  45. case "sqlite":
  46. default:
  47. _db = new SqliteDatabase();
  48. break;
  49. }
  50. _db.Configure(config);
  51. }
  52. /// <summary>
  53. /// 用默认 SQLite 创建产品数据库(按产品名建文件,放产品目录下)
  54. /// </summary>
  55. /// <param name="productDir">产品目录(如 Products\ProductA\)</param>
  56. /// <param name="productName">产品名</param>
  57. public static ProductDatabase CreateDefaultSqlite(string productDir, string productName)
  58. {
  59. var dbPath = Path.Combine(productDir, "product.db");
  60. var config = new DatabaseConfig
  61. {
  62. Id = Guid.NewGuid(),
  63. Name = productName,
  64. ProviderType = "sqlite",
  65. Server = dbPath
  66. };
  67. return new ProductDatabase(config);
  68. }
  69. /// <summary>
  70. /// 初始化:打开连接 + 自动建表
  71. /// </summary>
  72. public (bool Success, string Message) Initialize()
  73. {
  74. try
  75. {
  76. // SQLite 自动建目录
  77. if ((_config.ProviderType ?? "sqlite").ToLowerInvariant() == "sqlite")
  78. {
  79. var filePath = _config.Server;
  80. var dir = Path.GetDirectoryName(filePath);
  81. if (!string.IsNullOrWhiteSpace(dir) && !Directory.Exists(dir))
  82. Directory.CreateDirectory(dir);
  83. }
  84. if (!_db.Open())
  85. return (false, "产品数据库打开失败");
  86. CreateTables();
  87. return (true, "产品数据库初始化成功");
  88. }
  89. catch (Exception ex)
  90. {
  91. return (false, $"产品数据库初始化失败:{ex.Message}");
  92. }
  93. }
  94. #region 建表
  95. private void CreateTables()
  96. {
  97. var sqls = new List<string>
  98. {
  99. // 生产记录
  100. @"CREATE TABLE IF NOT EXISTS prod_records (
  101. id INTEGER PRIMARY KEY AUTOINCREMENT,
  102. recipe_name TEXT,
  103. batch_no TEXT,
  104. serial_no TEXT,
  105. result INTEGER DEFAULT 0,
  106. start_time TEXT,
  107. end_time TEXT,
  108. duration_ms INTEGER,
  109. operator_name TEXT,
  110. remark TEXT,
  111. created_at TEXT DEFAULT (datetime('now','localtime'))
  112. )",
  113. // 运行日志
  114. @"CREATE TABLE IF NOT EXISTS run_logs (
  115. id INTEGER PRIMARY KEY AUTOINCREMENT,
  116. log_level TEXT DEFAULT 'INFO',
  117. module TEXT,
  118. message TEXT,
  119. detail TEXT,
  120. created_at TEXT DEFAULT (datetime('now','localtime'))
  121. )",
  122. // 报警历史
  123. @"CREATE TABLE IF NOT EXISTS alarm_history (
  124. id INTEGER PRIMARY KEY AUTOINCREMENT,
  125. alarm_code TEXT,
  126. alarm_name TEXT,
  127. alarm_level INTEGER DEFAULT 2,
  128. description TEXT,
  129. source TEXT,
  130. start_time TEXT,
  131. end_time TEXT,
  132. ack_by TEXT,
  133. ack_time TEXT,
  134. is_acknowledged INTEGER DEFAULT 0,
  135. is_resolved INTEGER DEFAULT 0,
  136. created_at TEXT DEFAULT (datetime('now','localtime'))
  137. )",
  138. // 设备运行数据(时序型)
  139. @"CREATE TABLE IF NOT EXISTS device_data (
  140. id INTEGER PRIMARY KEY AUTOINCREMENT,
  141. device_name TEXT NOT NULL,
  142. data_key TEXT NOT NULL,
  143. data_value REAL,
  144. value_text TEXT,
  145. quality INTEGER DEFAULT 1,
  146. created_at TEXT DEFAULT (datetime('now','localtime'))
  147. )",
  148. // 产品信息表
  149. @"CREATE TABLE IF NOT EXISTS product_info (
  150. id INTEGER PRIMARY KEY AUTOINCREMENT,
  151. product_name TEXT NOT NULL UNIQUE,
  152. version TEXT,
  153. description TEXT,
  154. created_at TEXT DEFAULT (datetime('now','localtime')),
  155. updated_at TEXT DEFAULT (datetime('now','localtime'))
  156. )"
  157. };
  158. // 索引
  159. sqls.Add("CREATE INDEX IF NOT EXISTS idx_prod_result ON prod_records(result)");
  160. sqls.Add("CREATE INDEX IF NOT EXISTS idx_prod_time ON prod_records(start_time)");
  161. sqls.Add("CREATE INDEX IF NOT EXISTS idx_prod_serial ON prod_records(serial_no)");
  162. sqls.Add("CREATE INDEX IF NOT EXISTS idx_log_level ON run_logs(log_level)");
  163. sqls.Add("CREATE INDEX IF NOT EXISTS idx_log_time ON run_logs(created_at)");
  164. sqls.Add("CREATE INDEX IF NOT EXISTS idx_log_module ON run_logs(module)");
  165. sqls.Add("CREATE INDEX IF NOT EXISTS idx_alarm_level ON alarm_history(alarm_level)");
  166. sqls.Add("CREATE INDEX IF NOT EXISTS idx_alarm_time ON alarm_history(start_time)");
  167. sqls.Add("CREATE INDEX IF NOT EXISTS idx_alarm_resolved ON alarm_history(is_resolved)");
  168. sqls.Add("CREATE INDEX IF NOT EXISTS idx_device_name ON device_data(device_name)");
  169. sqls.Add("CREATE INDEX IF NOT EXISTS idx_device_time ON device_data(created_at)");
  170. sqls.Add("CREATE INDEX IF NOT EXISTS idx_device_key ON device_data(data_key)");
  171. _db.ExecuteTransaction(sqls);
  172. // 写入产品信息
  173. var productInfo = _db.ExecuteScalar(
  174. "SELECT COUNT(*) FROM product_info WHERE product_name = @n",
  175. new Dictionary<string, object> {{"@n", ProductName}});
  176. if (Convert.ToInt32(productInfo) == 0)
  177. {
  178. _db.ExecuteNonQuery(
  179. "INSERT INTO product_info (product_name, description) VALUES (@n, @d)",
  180. new Dictionary<string, object>
  181. {
  182. {"@n", ProductName},
  183. {"@d", $"产品[{ProductName}]的运行数据库"}
  184. });
  185. }
  186. }
  187. #endregion
  188. #region 生产记录
  189. /// <summary>
  190. /// 写入一条生产记录(开始)
  191. /// </summary>
  192. /// <returns>记录ID</returns>
  193. public int StartProduction(string recipeName, string batchNo, string serialNo, string operatorName = null)
  194. {
  195. var sql = @"INSERT INTO prod_records
  196. (recipe_name, batch_no, serial_no, start_time, operator_name)
  197. VALUES (@r, @b, @s, datetime('now','localtime'), @o)";
  198. var p = new Dictionary<string, object>
  199. {
  200. {"@r", (object)recipeName ?? DBNull.Value},
  201. {"@b", (object)batchNo ?? DBNull.Value},
  202. {"@s", (object)serialNo ?? DBNull.Value},
  203. {"@o", (object)operatorName ?? DBNull.Value}
  204. };
  205. _db.ExecuteNonQuery(sql, p);
  206. return Convert.ToInt32(_db.ExecuteScalar("SELECT last_insert_rowid()"));
  207. }
  208. /// <summary>
  209. /// 完成生产记录
  210. /// </summary>
  211. public void FinishProduction(int recordId, bool success, string remark = null)
  212. {
  213. var sql = @"UPDATE prod_records
  214. SET result = @r, end_time = datetime('now','localtime'),
  215. duration_ms = CAST((julianday('now','localtime') - julianday(start_time)) * 86400000 AS INTEGER),
  216. remark = @rk
  217. WHERE id = @id";
  218. var p = new Dictionary<string, object>
  219. {
  220. {"@id", recordId},
  221. {"@r", success ? 1 : 0},
  222. {"@rk", (object)remark ?? DBNull.Value}
  223. };
  224. _db.ExecuteNonQuery(sql, p);
  225. }
  226. /// <summary>
  227. /// 查询生产记录统计
  228. /// </summary>
  229. public (int Total, int Ok, int Ng, double Yield) GetProductionStats(DateTime? startTime = null, DateTime? endTime = null)
  230. {
  231. var conditions = new List<string>();
  232. var parameters = new Dictionary<string, object>();
  233. if (startTime.HasValue)
  234. {
  235. conditions.Add("start_time >= @st");
  236. parameters["@st"] = startTime.Value.ToString("yyyy-MM-dd HH:mm:ss");
  237. }
  238. if (endTime.HasValue)
  239. {
  240. conditions.Add("start_time <= @et");
  241. parameters["@et"] = endTime.Value.ToString("yyyy-MM-dd HH:mm:ss");
  242. }
  243. var where = conditions.Count > 0 ? "WHERE " + string.Join(" AND ", conditions) : "";
  244. var total = Convert.ToInt32(
  245. _db.ExecuteScalar($"SELECT COUNT(*) FROM prod_records {where}", parameters));
  246. var ok = Convert.ToInt32(
  247. _db.ExecuteScalar($"SELECT COUNT(*) FROM prod_records {where} AND result = 1",
  248. parameters));
  249. var ng = total - ok;
  250. var yield = total > 0 ? (double)ok / total * 100 : 0;
  251. return (total, ok, ng, Math.Round(yield, 2));
  252. }
  253. #endregion
  254. #region 运行日志
  255. /// <summary>
  256. /// 写入运行日志
  257. /// </summary>
  258. public void WriteLog(string level, string module, string message, string detail = null)
  259. {
  260. try
  261. {
  262. var sql = @"INSERT INTO run_logs (log_level, module, message, detail)
  263. VALUES (@l, @m, @msg, @d)";
  264. var p = new Dictionary<string, object>
  265. {
  266. {"@l", level ?? "INFO"},
  267. {"@m", (object)module ?? DBNull.Value},
  268. {"@msg", message ?? ""},
  269. {"@d", (object)detail ?? DBNull.Value}
  270. };
  271. _db.ExecuteNonQuery(sql, p);
  272. }
  273. catch { /* 日志写入失败不抛 */ }
  274. }
  275. /// <summary>
  276. /// 便捷方法:INFO
  277. /// </summary>
  278. public void LogInfo(string module, string message, string detail = null)
  279. => WriteLog("INFO", module, message, detail);
  280. /// <summary>
  281. /// 便捷方法:WARN
  282. /// </summary>
  283. public void LogWarn(string module, string message, string detail = null)
  284. => WriteLog("WARN", module, message, detail);
  285. /// <summary>
  286. /// 便捷方法:ERROR
  287. /// </summary>
  288. public void LogError(string module, string message, string detail = null)
  289. => WriteLog("ERROR", module, message, detail);
  290. #endregion
  291. #region 报警
  292. /// <summary>
  293. /// 触发报警
  294. /// </summary>
  295. /// <returns>报警记录ID</returns>
  296. public int RaiseAlarm(string code, string name, int level, string description = null, string source = null)
  297. {
  298. var sql = @"INSERT INTO alarm_history
  299. (alarm_code, alarm_name, alarm_level, description, source, start_time)
  300. VALUES (@c, @n, @l, @d, @s, datetime('now','localtime'))";
  301. var p = new Dictionary<string, object>
  302. {
  303. {"@c", code},
  304. {"@n", name},
  305. {"@l", level},
  306. {"@d", (object)description ?? DBNull.Value},
  307. {"@s", (object)source ?? DBNull.Value}
  308. };
  309. _db.ExecuteNonQuery(sql, p);
  310. return Convert.ToInt32(_db.ExecuteScalar("SELECT last_insert_rowid()"));
  311. }
  312. /// <summary>
  313. /// 确认报警
  314. /// </summary>
  315. public void AcknowledgeAlarm(int alarmId, string ackBy = null)
  316. {
  317. var sql = @"UPDATE alarm_history
  318. SET is_acknowledged = 1, ack_by = @a, ack_time = datetime('now','localtime')
  319. WHERE id = @id";
  320. var p = new Dictionary<string, object>
  321. {
  322. {"@id", alarmId},
  323. {"@a", (object)ackBy ?? DBNull.Value}
  324. };
  325. _db.ExecuteNonQuery(sql, p);
  326. }
  327. /// <summary>
  328. /// 消除报警
  329. /// </summary>
  330. public void ResolveAlarm(int alarmId)
  331. {
  332. var sql = @"UPDATE alarm_history
  333. SET is_resolved = 1, end_time = datetime('now','localtime')
  334. WHERE id = @id";
  335. _db.ExecuteNonQuery(sql, new Dictionary<string, object> {{"@id", alarmId}});
  336. }
  337. /// <summary>
  338. /// 获取未消除报警数
  339. /// </summary>
  340. public int GetActiveAlarmCount()
  341. {
  342. return Convert.ToInt32(
  343. _db.ExecuteScalar("SELECT COUNT(*) FROM alarm_history WHERE is_resolved = 0"));
  344. }
  345. #endregion
  346. #region 设备数据
  347. /// <summary>
  348. /// 写入设备运行数据
  349. /// </summary>
  350. public void WriteDeviceData(string deviceName, string dataKey, double? value = null, string valueText = null, int quality = 1)
  351. {
  352. try
  353. {
  354. var sql = @"INSERT INTO device_data (device_name, data_key, data_value, value_text, quality)
  355. VALUES (@d, @k, @v, @vt, @q)";
  356. var p = new Dictionary<string, object>
  357. {
  358. {"@d", deviceName},
  359. {"@k", dataKey},
  360. {"@v", (object)value ?? DBNull.Value},
  361. {"@vt", (object)valueText ?? DBNull.Value},
  362. {"@q", quality}
  363. };
  364. _db.ExecuteNonQuery(sql, p);
  365. }
  366. catch { /* 数据写入失败不抛 */ }
  367. }
  368. /// <summary>
  369. /// 批量写入设备数据
  370. /// </summary>
  371. public void WriteDeviceDataBatch(IEnumerable<(string deviceName, string key, double? value, string valueText)> dataList)
  372. {
  373. if (dataList == null) return;
  374. var sqls = new List<string>();
  375. // 简单拼装(内部数据,安全风险低)
  376. foreach (var d in dataList)
  377. {
  378. var v = d.value.HasValue
  379. ? d.value.Value.ToString(System.Globalization.CultureInfo.InvariantCulture)
  380. : "NULL";
  381. var vt = d.valueText != null ? $"'{d.valueText.Replace("'", "''")}'" : "NULL";
  382. sqls.Add($@"INSERT INTO device_data (device_name, data_key, data_value, value_text)
  383. VALUES ('{d.deviceName.Replace("'", "''")}', '{d.key.Replace("'", "''")}', {v}, {vt})");
  384. }
  385. if (sqls.Count > 0)
  386. _db.ExecuteTransaction(sqls);
  387. }
  388. /// <summary>
  389. /// 查询设备某数据点的历史
  390. /// </summary>
  391. public DataTable QueryDeviceHistory(string deviceName, string dataKey,
  392. DateTime? startTime = null, DateTime? endTime = null, int limit = 1000)
  393. {
  394. var conditions = new List<string> {"device_name = @d", "data_key = @k"};
  395. var parameters = new Dictionary<string, object>
  396. {
  397. {"@d", deviceName},
  398. {"@k", dataKey}
  399. };
  400. if (startTime.HasValue)
  401. {
  402. conditions.Add("created_at >= @st");
  403. parameters["@st"] = startTime.Value.ToString("yyyy-MM-dd HH:mm:ss");
  404. }
  405. if (endTime.HasValue)
  406. {
  407. conditions.Add("created_at <= @et");
  408. parameters["@et"] = endTime.Value.ToString("yyyy-MM-dd HH:mm:ss");
  409. }
  410. var where = "WHERE " + string.Join(" AND ", conditions);
  411. var sql = $@"
  412. SELECT * FROM device_data {where}
  413. ORDER BY id DESC
  414. LIMIT @lim";
  415. parameters["@lim"] = limit;
  416. var dt = _db.ExecuteQuery(sql, parameters);
  417. // 反序变成时间正序
  418. var result = dt.Clone();
  419. for (int i = dt.Rows.Count - 1; i >= 0; i--)
  420. {
  421. result.ImportRow(dt.Rows[i]);
  422. }
  423. return result;
  424. }
  425. #endregion
  426. #region IDisposable
  427. private bool _disposed;
  428. public void Dispose()
  429. {
  430. Dispose(true);
  431. GC.SuppressFinalize(this);
  432. }
  433. protected virtual void Dispose(bool disposing)
  434. {
  435. if (_disposed) return;
  436. if (disposing)
  437. {
  438. try { _db?.Dispose(); } catch { }
  439. }
  440. _disposed = true;
  441. }
  442. ~ProductDatabase()
  443. {
  444. Dispose(false);
  445. }
  446. #endregion
  447. }
  448. }