using System.Security.Cryptography; using System.Text.Json; using Admin.NET.Plugin.AiDOP.DataPlatform.Executors; using Admin.NET.Plugin.AiDOP.Entity.DataPlatform; using Microsoft.AspNetCore.Http; using MiniExcelLibs; namespace Admin.NET.Plugin.AiDOP.DataPlatform.FileImport; public sealed class MdpExcelImportService : ITransient { private readonly ISqlSugarClient _db; private readonly UserManager _userManager; private readonly MdpStagingWriter _writer; private readonly MdpExcelTemplateService _template; private readonly IEnumerable _consumers; public MdpExcelImportService( ISqlSugarClient db, UserManager userManager, MdpStagingWriter writer, MdpExcelTemplateService template, IEnumerable consumers) { _db = db; _userManager = userManager; _writer = writer; _template = template; _consumers = consumers; } private long TenantId => AidopTenantScope.ResolveOrThrow(_userManager); public async Task ListEntitiesAsync(CancellationToken ct) { var tenantId = TenantId; var sources = await _db.Queryable() .Where(s => s.SourceType == "FILE_EXCEL" && s.Status == 1) .ToListAsync(ct); var sourceIds = sources.Select(s => s.Id).ToList(); var entities = await _db.Queryable() .Where(e => sourceIds.Contains(e.SourceId) && e.Status == 1 && (e.TenantId == tenantId || e.IsPublicTemplate == 1)) .ToListAsync(ct); return new { list = entities.Select(e => { var s = sources.First(x => x.Id == e.SourceId); return new { e.Id, e.EntityCode, e.EntityName, s.SourceCode, e.TargetTableName, e.BizKeyExpr, e.IsPublicTemplate }; }) }; } public async Task TemplateAsync(string entityCode, CancellationToken ct) { var (entity, _, mappings) = await LoadEntityAsync(entityCode, ct); return await _template.BuildAsync(entity, mappings); } public async Task ValidateAsync(IFormFile file, string entityCode, long factoryId, CancellationToken ct) { var tenantId = TenantId; var (entity, source, mappings) = await LoadEntityAsync(entityCode, ct); if (file == null || file.Length == 0) throw Oops.Oh("请上传 Excel 文件"); await using var stream = file.OpenReadStream(); MdpFileImportSecurity.EnsureXlsx(file.FileName, file.Length, stream); var sha = await ShaAsync(stream, ct); stream.Position = 0; var dup = await _db.Queryable() .Where(x => x.TenantId == tenantId && x.FileSha256 == sha && x.Status == "SUCCESS") .CountAsync(ct); if (dup > 0) throw Oops.Oh("相同文件已成功导入"); var sheet = MdpExcelTemplateService.ReadSheet(entity); var rows = MiniExcel.Query(stream, useHeaderRow: true, sheetName: sheet).Cast>().ToList(); if (rows.Count > MdpFileImportSecurity.MaxRows) throw Oops.Oh($"超过最大行数 {MdpFileImportSecurity.MaxRows}"); var now = DateTime.Now; var batch = new MdpFileImportBatch { BatchNo = $"MDP_FILE_{now:yyyyMMddHHmmss}_{Random.Shared.Next(1000, 9999)}", TenantId = tenantId, FactoryId = factoryId, SourceId = source.Id, EntityId = entity.Id, EntityCode = entity.EntityCode, FileName = file.FileName, FileSha256 = sha, FileSize = file.Length, SheetName = sheet, ImportMode = "UPSERT", Status = "VALIDATING", CreatedBy = _userManager.RealName ?? _userManager.Account, CreateTime = now }; batch.Id = await _db.Insertable(batch).ExecuteReturnIdentityAsync(); var errors = new List(); var hashes = new HashSet(StringComparer.OrdinalIgnoreCase); var valid = 0; for (var i = 0; i < rows.Count; i++) { var rowNo = i + 2; try { var raw = rows[i].ToDictionary(x => StripStar(x.Key), x => x.Value, StringComparer.OrdinalIgnoreCase); MdpFileImportSecurity.StripClientTenant(raw!); var mapped = MdpExcelValidation.MapRow(raw!, mappings); mapped["factory_id"] = factoryId; var biz = MdpExcelValidation.BusinessKey(entity.BizKeyExpr, mapped); var hash = MdpExcelValidation.RowHash(mapped); var writeStatus = hashes.Add(hash) ? "PENDING" : "DUPLICATE"; if (writeStatus == "DUPLICATE") errors.Add(Err(batch.Id, rowNo, "ROW_HASH", "DUPLICATE_HASH", "相同内容重复")); else valid++; await _db.Insertable(new MdpFileImportRow { BatchId = batch.Id, RowNo = rowNo, BusinessKey = biz, RowHash = hash, RawJson = JsonSerializer.Serialize(raw), MappedJson = JsonSerializer.Serialize(mapped), ValidationStatus = writeStatus == "PENDING" ? "OK" : "ERROR", WriteStatus = writeStatus, ErrorCount = writeStatus == "PENDING" ? 0 : 1 }).ExecuteCommandAsync(); } catch (Exception ex) { errors.Add(Err(batch.Id, rowNo, null, "ROW_INVALID", ex.Message)); await _db.Insertable(new MdpFileImportRow { BatchId = batch.Id, RowNo = rowNo, ValidationStatus = "ERROR", WriteStatus = "REJECTED", ErrorCount = 1, RawJson = JsonSerializer.Serialize(rows[i]) }).ExecuteCommandAsync(); } } if (errors.Count > 0) await _db.Insertable(errors).ExecuteCommandAsync(); batch.TotalRows = rows.Count; batch.ValidRows = valid; batch.ErrorRows = errors.Count; batch.Status = errors.Count == 0 && valid > 0 ? "READY" : "VALIDATION_FAILED"; await _db.Updateable(batch).ExecuteCommandAsync(); return new { batchNo = batch.BatchNo, batch.Status, batch.TotalRows, batch.ValidRows, batch.ErrorRows }; } public async Task GetAsync(string batchNo, CancellationToken ct) { var tenantId = TenantId; var batch = await _db.Queryable().FirstAsync(x => x.TenantId == tenantId && x.BatchNo == batchNo, ct) ?? throw Oops.Oh("批次不存在"); return batch; } public async Task ErrorsAsync(string batchNo, CancellationToken ct) { var tenantId = TenantId; var batch = await _db.Queryable().FirstAsync(x => x.TenantId == tenantId && x.BatchNo == batchNo, ct) ?? throw Oops.Oh("批次不存在"); var errors = await _db.Queryable().Where(x => x.BatchId == batch.Id).ToListAsync(ct); var memory = new MemoryStream(); await memory.SaveAsAsync(errors, sheetName: "错误"); memory.Seek(0, SeekOrigin.Begin); return new FileStreamResult(memory, "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet") { FileDownloadName = $"{batch.BatchNo}-errors.xlsx" }; } public async Task CommitAsync(string batchNo, CancellationToken ct) { var tenantId = TenantId; var batch = await _db.Queryable().FirstAsync(x => x.TenantId == tenantId && x.BatchNo == batchNo, ct) ?? throw Oops.Oh("批次不存在"); if (batch.Status == "SUCCESS") throw Oops.Oh("该批次已导入"); if (batch.Status != "READY") throw Oops.Oh("只有预检通过的批次可以确认导入"); var (entity, source, _) = await LoadEntityAsync(batch.EntityCode, ct); if (string.IsNullOrWhiteSpace(entity.TargetTableName) || !System.Text.RegularExpressions.Regex.IsMatch(entity.TargetTableName, @"^[A-Za-z0-9_]+$")) throw Oops.Oh("实体目标表非法"); batch.Status = "IMPORTING"; await _db.Updateable(batch).UpdateColumns(x => new { x.Status }).ExecuteCommandAsync(); var rows = await _db.Queryable() .Where(x => x.BatchId == batch.Id && x.ValidationStatus == "OK") .ToListAsync(ct); var ctx = new MdpPullContext { TenantId = tenantId, BatchId = batch.BatchNo }; var written = 0; foreach (var row in rows) { var mapped = JsonSerializer.Deserialize>(row.MappedJson ?? "{}") ?? new Dictionary(); MdpFileImportSecurity.StripClientTenant(mapped); mapped["factory_id"] = batch.FactoryId; var n = await _writer.UpsertAsync( source, entity, entity.EntityCode, mapped, row.RawJson ?? "{}", $"{batch.BatchNo}:{row.RowNo}", ctx); written += n; row.WriteStatus = n > 0 ? "WRITTEN" : "SKIPPED"; await _db.Updateable(row).UpdateColumns(x => new { x.WriteStatus }).ExecuteCommandAsync(); } var consumer = _consumers.FirstOrDefault(c => c.CanConsume(entity.EntityCode)); if (consumer != null) { await consumer.ConsumeAsync(entity.EntityCode, tenantId, batch.FactoryId, batch.BatchNo, ct); batch.TransformStatus = "SUCCESS"; batch.TransformTime = DateTime.Now; } batch.WrittenRows = written; batch.Status = "SUCCESS"; batch.CommitTime = DateTime.Now; await _db.Updateable(batch).ExecuteCommandAsync(); return new { batch.BatchNo, batch.Status, batch.WrittenRows, batch.TransformStatus }; } public async Task VoidAsync(string batchNo, CancellationToken ct) { var tenantId = TenantId; var batch = await _db.Queryable().FirstAsync(x => x.TenantId == tenantId && x.BatchNo == batchNo, ct) ?? throw Oops.Oh("批次不存在"); if (batch.Status is "SUCCESS" or "IMPORTING") throw Oops.Oh("已提交批次请使用回滚"); batch.Status = "VOIDED"; await _db.Updateable(batch).UpdateColumns(x => new { x.Status }).ExecuteCommandAsync(); return new { ok = true }; } public async Task RollbackAsync(string batchNo, CancellationToken ct) { var tenantId = TenantId; var batch = await _db.Queryable().FirstAsync(x => x.TenantId == tenantId && x.BatchNo == batchNo, ct) ?? throw Oops.Oh("批次不存在"); if (batch.Status != "SUCCESS") throw Oops.Oh("仅成功批次可回滚"); var entity = await _db.Queryable().FirstAsync(x => x.Id == batch.EntityId, ct) ?? throw Oops.Oh("实体不存在"); var rebuildable = entity.FileConfigJson?.Contains("\"rebuildable\":true", StringComparison.OrdinalIgnoreCase) == true; if (!string.IsNullOrWhiteSpace(batch.TransformStatus) && batch.TransformStatus != "PENDING" && !rebuildable) throw Oops.Oh("已 Transform 且实体不可重建,请业务冲销"); if (string.IsNullOrWhiteSpace(entity.TargetTableName) || !System.Text.RegularExpressions.Regex.IsMatch(entity.TargetTableName, @"^[A-Za-z0-9_]+$")) throw Oops.Oh("实体目标表非法"); await _db.Ado.ExecuteCommandAsync( $"DELETE FROM `{entity.TargetTableName}` WHERE tenant_id=@TenantId AND sync_batch_id=@Batch", new SugarParameter("@TenantId", tenantId), new SugarParameter("@Batch", batch.BatchNo)); batch.Status = "ROLLED_BACK"; await _db.Updateable(batch).UpdateColumns(x => new { x.Status }).ExecuteCommandAsync(); return new { ok = true }; } private async Task<(MdpEntity entity, MdpSource source, List mappings)> LoadEntityAsync(string entityCode, CancellationToken ct) { var tenantId = TenantId; var entity = await _db.Queryable() .Where(x => x.EntityCode == entityCode && x.Status == 1 && (x.TenantId == tenantId || x.IsPublicTemplate == 1)) .FirstAsync(ct) ?? throw Oops.Oh("FILE 实体不存在或无权访问"); var source = await _db.Queryable().FirstAsync(x => x.Id == entity.SourceId, ct) ?? throw Oops.Oh("数据源不存在"); if (!string.Equals(source.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase)) throw Oops.Oh("该实体不是 FILE_EXCEL"); var mappings = await _db.Queryable().Where(x => x.EntityId == entity.Id).OrderBy(x => x.SortOrder).ToListAsync(ct); return (entity, source, mappings); } private static string StripStar(string name) => name.Replace(" ★", string.Empty).Trim(); private static MdpFileImportError Err(long batchId, int rowNo, string? field, string code, string message) => new() { BatchId = batchId, RowNo = rowNo, FieldName = field, ErrorCode = code, ErrorMessage = message }; private static async Task ShaAsync(Stream stream, CancellationToken ct) { using var sha = SHA256.Create(); var hash = await sha.ComputeHashAsync(stream, ct); return Convert.ToHexString(hash).ToLowerInvariant(); } } public interface IMdpFileImportConsumer { bool CanConsume(string entityCode); Task ConsumeAsync(string entityCode, long tenantId, long factoryId, string batchNo, CancellationToken ct); }