MdpSyncEntityConfigService.cs 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391
  1. using Admin.NET.Plugin.AiDOP.Dto.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  3. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  4. /// <summary>
  5. /// 同步任务实体与字段映射配置 API。
  6. /// </summary>
  7. [ApiDescriptionSettings(Order = 324, Description = "数据中台同步实体与字段映射")]
  8. [Route("api/DataPlatform")]
  9. [NonUnify]
  10. public class MdpSyncEntityConfigService : IDynamicApiController, ITransient
  11. {
  12. private static readonly HashSet<string> ValidFieldTypes = new(StringComparer.OrdinalIgnoreCase)
  13. {
  14. "DIRECT", "JSONPATH", "SCRIPT", "CONST", "LOOKUP"
  15. };
  16. private readonly ISqlSugarClient _db;
  17. private readonly UserManager _userManager;
  18. public MdpSyncEntityConfigService(ISqlSugarClient db, UserManager userManager)
  19. {
  20. _db = db;
  21. _userManager = userManager;
  22. }
  23. [DisplayName("同步任务实体列表")]
  24. [HttpGet("sync-tasks/{taskCode}/entities")]
  25. public async Task<object> GetEntities(string taskCode)
  26. {
  27. var tenantId = _userManager.TenantId;
  28. var task = await FindTaskAsync(taskCode, tenantId);
  29. var prefix = ResolveEntityCodePrefix(task.TaskCode);
  30. var jobCode = string.IsNullOrWhiteSpace(task.JobCode) ? task.TaskCode : task.JobCode;
  31. var where = new List<string> { "(e.tenant_id = @TenantId OR e.tenant_id = 0)", "e.status = 1" };
  32. var pars = new List<SugarParameter> { new("@TenantId", tenantId) };
  33. if (!string.IsNullOrWhiteSpace(prefix))
  34. {
  35. where.Add("(e.entity_code LIKE @EntityPrefix OR e.job_id = @JobCode)");
  36. pars.Add(new SugarParameter("@EntityPrefix", prefix + "%"));
  37. pars.Add(new SugarParameter("@JobCode", jobCode));
  38. }
  39. else
  40. {
  41. where.Add("e.job_id = @JobCode");
  42. pars.Add(new SugarParameter("@JobCode", jobCode));
  43. }
  44. var whereSql = string.Join(" AND ", where);
  45. var list = await _db.Ado.SqlQueryAsync<MdpEntityRow>(
  46. $"""
  47. SELECT e.id AS Id, e.tenant_id AS TenantId, e.source_id AS SourceId,
  48. s.source_code AS SourceCode, s.source_name AS SourceName, s.source_type AS SourceType,
  49. e.entity_code AS EntityCode, e.entity_name AS EntityName, e.entity_type AS EntityType,
  50. e.source_table_name AS SourceTableName, e.source_api_path AS SourceApiPath,
  51. e.target_table_name AS TargetTableName, e.sync_mode AS SyncMode, e.incr_column AS IncrColumn,
  52. e.batch_size AS BatchSize, e.job_id AS JobId, e.status AS Status, e.remark AS Remark,
  53. IFNULL(m.cnt, 0) AS FieldMappingCount
  54. FROM mdp_entity e
  55. LEFT JOIN mdp_source s ON s.id = e.source_id
  56. LEFT JOIN (
  57. SELECT entity_id, COUNT(1) AS cnt FROM mdp_field_mapping GROUP BY entity_id
  58. ) m ON m.entity_id = e.id
  59. WHERE {whereSql}
  60. ORDER BY e.entity_code
  61. """,
  62. pars);
  63. return list;
  64. }
  65. [DisplayName("新增同步任务实体")]
  66. [HttpPost("sync-tasks/{taskCode}/entities")]
  67. public async Task<object> CreateEntity(string taskCode, [FromBody] MdpEntityUpsertInput input)
  68. {
  69. ValidateEntityUpsert(input);
  70. var tenantId = _userManager.TenantId;
  71. var task = await FindTaskAsync(taskCode, tenantId);
  72. var sourceId = await ResolveSourceIdAsync(input, tenantId);
  73. var entityTenantId = task.TenantId == 0 ? 0L : tenantId;
  74. var jobCode = string.IsNullOrWhiteSpace(task.JobCode) ? task.TaskCode : task.JobCode;
  75. var exists = await _db.Queryable<MdpEntity>()
  76. .Where(u => u.TenantId == entityTenantId && u.EntityCode == input.EntityCode.Trim())
  77. .AnyAsync();
  78. if (exists)
  79. throw Oops.Oh("实体编码已存在");
  80. var now = DateTime.Now;
  81. var entity = new MdpEntity
  82. {
  83. TenantId = entityTenantId,
  84. SourceId = sourceId,
  85. EntityCode = input.EntityCode.Trim(),
  86. EntityName = input.EntityName.Trim(),
  87. EntityType = string.IsNullOrWhiteSpace(input.EntityType) ? "TABLE" : input.EntityType.Trim(),
  88. SourceTableName = input.SourceTableName?.Trim(),
  89. SourceApiPath = input.SourceApiPath?.Trim(),
  90. TargetTableName = input.TargetTableName?.Trim(),
  91. SyncMode = string.IsNullOrWhiteSpace(input.SyncMode) ? "INCR" : input.SyncMode.Trim(),
  92. IncrColumn = input.IncrColumn?.Trim(),
  93. BatchSize = input.BatchSize <= 0 ? 5000 : input.BatchSize,
  94. JobId = string.IsNullOrWhiteSpace(input.JobId) ? jobCode : input.JobId.Trim(),
  95. Status = input.Status <= 0 ? 0 : 1,
  96. Remark = input.Remark?.Trim(),
  97. InboundEnabled = await ResolveInboundEnabledAsync(sourceId, input),
  98. BizKeyExpr = input.BizKeyExpr?.Trim(),
  99. CreateTime = now,
  100. UpdateTime = now
  101. };
  102. var id = await _db.Insertable(entity).ExecuteReturnBigIdentityAsync();
  103. await BumpTaskConfigVersionAsync(task);
  104. return new { id };
  105. }
  106. [DisplayName("更新同步实体")]
  107. [HttpPut("entities/{id:long}")]
  108. public async Task<object> UpdateEntity(long id, [FromBody] MdpEntityUpsertInput input)
  109. {
  110. ValidateEntityUpsert(input);
  111. var tenantId = _userManager.TenantId;
  112. var entity = await FindEditableEntityAsync(id, tenantId);
  113. entity.EntityName = input.EntityName.Trim();
  114. entity.EntityType = string.IsNullOrWhiteSpace(input.EntityType) ? "TABLE" : input.EntityType.Trim();
  115. entity.SourceId = await ResolveSourceIdAsync(input, tenantId, entity.SourceId);
  116. entity.SourceTableName = input.SourceTableName?.Trim();
  117. entity.SourceApiPath = input.SourceApiPath?.Trim();
  118. entity.TargetTableName = input.TargetTableName?.Trim();
  119. entity.SyncMode = string.IsNullOrWhiteSpace(input.SyncMode) ? "INCR" : input.SyncMode.Trim();
  120. entity.IncrColumn = input.IncrColumn?.Trim();
  121. entity.BatchSize = input.BatchSize <= 0 ? 5000 : input.BatchSize;
  122. entity.JobId = input.JobId?.Trim();
  123. entity.Status = input.Status <= 0 ? 0 : 1;
  124. entity.Remark = input.Remark?.Trim();
  125. entity.InboundEnabled = await ResolveInboundEnabledAsync(entity.SourceId, input);
  126. entity.BizKeyExpr = input.BizKeyExpr?.Trim();
  127. entity.UpdateTime = DateTime.Now;
  128. if (!IsBuiltInEntity(entity))
  129. entity.EntityCode = input.EntityCode.Trim();
  130. await _db.Updateable(entity).ExecuteCommandAsync();
  131. return new { id = entity.Id };
  132. }
  133. [DisplayName("删除同步实体")]
  134. [HttpDelete("entities/{id:long}")]
  135. public async Task<object> DeleteEntity(long id)
  136. {
  137. var tenantId = _userManager.TenantId;
  138. var entity = await FindEditableEntityAsync(id, tenantId);
  139. if (IsBuiltInEntity(entity) && entity.TenantId == 0)
  140. throw Oops.Oh("内置 MDP 实体不允许删除");
  141. entity.Status = 0;
  142. entity.UpdateTime = DateTime.Now;
  143. await _db.Updateable(entity).UpdateColumns(u => new { u.Status, u.UpdateTime }).ExecuteCommandAsync();
  144. return new { id = entity.Id };
  145. }
  146. [DisplayName("实体字段映射列表")]
  147. [HttpGet("entities/{entityId:long}/field-mappings")]
  148. public async Task<object> GetFieldMappings(long entityId)
  149. {
  150. var tenantId = _userManager.TenantId;
  151. await FindEditableEntityAsync(entityId, tenantId);
  152. var list = await _db.Queryable<MdpFieldMapping>()
  153. .Where(u => u.EntityId == entityId)
  154. .OrderBy(u => u.SortOrder)
  155. .OrderBy(u => u.TargetField)
  156. .ToListAsync();
  157. return list.Select(MapFieldMappingRow).ToList();
  158. }
  159. [DisplayName("新增实体字段映射")]
  160. [HttpPost("entities/{entityId:long}/field-mappings")]
  161. public async Task<object> CreateFieldMapping(long entityId, [FromBody] MdpFieldMappingUpsertInput input)
  162. {
  163. ValidateFieldMappingUpsert(input);
  164. var tenantId = _userManager.TenantId;
  165. await FindEditableEntityAsync(entityId, tenantId);
  166. var exists = await _db.Queryable<MdpFieldMapping>()
  167. .Where(u => u.EntityId == entityId && u.TargetField == input.TargetField.Trim())
  168. .AnyAsync();
  169. if (exists)
  170. throw Oops.Oh("目标字段映射已存在");
  171. var entity = MapFieldMappingInsert(entityId, input);
  172. var id = await _db.Insertable(entity).ExecuteReturnBigIdentityAsync();
  173. return new { id };
  174. }
  175. [DisplayName("更新字段映射")]
  176. [HttpPut("field-mappings/{id:long}")]
  177. public async Task<object> UpdateFieldMapping(long id, [FromBody] MdpFieldMappingUpsertInput input)
  178. {
  179. ValidateFieldMappingUpsert(input);
  180. var tenantId = _userManager.TenantId;
  181. var mapping = await _db.Queryable<MdpFieldMapping>().Where(u => u.Id == id).FirstAsync();
  182. if (mapping == null)
  183. throw Oops.Oh("字段映射不存在");
  184. await FindEditableEntityAsync(mapping.EntityId, tenantId);
  185. var duplicate = await _db.Queryable<MdpFieldMapping>()
  186. .Where(u => u.EntityId == mapping.EntityId && u.TargetField == input.TargetField.Trim() && u.Id != id)
  187. .AnyAsync();
  188. if (duplicate)
  189. throw Oops.Oh("目标字段映射已存在");
  190. mapping.SourceField = input.SourceField.Trim();
  191. mapping.TargetField = input.TargetField.Trim();
  192. mapping.FieldType = NormalizeFieldType(input.FieldType);
  193. mapping.TransformScript = input.TransformScript?.Trim();
  194. mapping.ConstValue = input.ConstValue?.Trim();
  195. mapping.LookupTable = input.LookupTable?.Trim();
  196. mapping.IsRequired = input.IsRequired ? 1 : 0;
  197. mapping.DefaultValue = input.DefaultValue?.Trim();
  198. mapping.SortOrder = input.SortOrder;
  199. await _db.Updateable(mapping).ExecuteCommandAsync();
  200. return new { id = mapping.Id };
  201. }
  202. [DisplayName("删除字段映射")]
  203. [HttpDelete("field-mappings/{id:long}")]
  204. public async Task<object> DeleteFieldMapping(long id)
  205. {
  206. var tenantId = _userManager.TenantId;
  207. var mapping = await _db.Queryable<MdpFieldMapping>().Where(u => u.Id == id).FirstAsync();
  208. if (mapping == null)
  209. throw Oops.Oh("字段映射不存在");
  210. await FindEditableEntityAsync(mapping.EntityId, tenantId);
  211. await _db.Deleteable<MdpFieldMapping>().Where(u => u.Id == id).ExecuteCommandAsync();
  212. return new { id };
  213. }
  214. private async Task<MdpSyncTask> FindTaskAsync(string taskCode, long tenantId)
  215. {
  216. if (string.IsNullOrWhiteSpace(taskCode))
  217. throw Oops.Oh("任务编码不能为空");
  218. var task = await _db.Queryable<MdpSyncTask>()
  219. .Where(u => u.TaskCode == taskCode && (u.TenantId == tenantId || u.TenantId == 0) && u.Status == 1)
  220. .OrderBy(u => u.TenantId, OrderByType.Desc)
  221. .FirstAsync();
  222. if (task == null)
  223. throw Oops.Oh("同步任务不存在");
  224. return task;
  225. }
  226. private async Task<MdpEntity> FindEditableEntityAsync(long id, long tenantId)
  227. {
  228. var entity = await _db.Queryable<MdpEntity>()
  229. .Where(u => u.Id == id && (u.TenantId == tenantId || u.TenantId == 0))
  230. .FirstAsync();
  231. if (entity == null)
  232. throw Oops.Oh("同步实体不存在");
  233. return entity;
  234. }
  235. private async Task<long> ResolveSourceIdAsync(MdpEntityUpsertInput input, long tenantId, long? fallbackSourceId = null)
  236. {
  237. if (input.SourceId > 0)
  238. return input.SourceId;
  239. if (!string.IsNullOrWhiteSpace(input.SourceCode))
  240. {
  241. var sourceId = await _db.Ado.GetLongAsync(
  242. """
  243. SELECT id FROM mdp_source
  244. WHERE source_code = @SourceCode AND (tenant_id = @TenantId OR tenant_id = 0)
  245. ORDER BY tenant_id DESC
  246. LIMIT 1
  247. """,
  248. new List<SugarParameter>
  249. {
  250. new("@SourceCode", input.SourceCode.Trim()),
  251. new("@TenantId", tenantId)
  252. });
  253. if (sourceId > 0)
  254. return sourceId;
  255. throw Oops.Oh($"数据源 {input.SourceCode} 不存在");
  256. }
  257. if (fallbackSourceId is > 0)
  258. return fallbackSourceId.Value;
  259. throw Oops.Oh("请指定 sourceId 或 sourceCode");
  260. }
  261. private async Task BumpTaskConfigVersionAsync(MdpSyncTask task)
  262. {
  263. task.ConfigVersion += 1;
  264. task.UpdateTime = DateTime.Now;
  265. await _db.Updateable(task).UpdateColumns(u => new { u.ConfigVersion, u.UpdateTime }).ExecuteCommandAsync();
  266. }
  267. private static string? ResolveEntityCodePrefix(string taskCode)
  268. {
  269. if (string.IsNullOrWhiteSpace(taskCode))
  270. return null;
  271. var first = taskCode.Split('_', StringSplitOptions.RemoveEmptyEntries).FirstOrDefault();
  272. if (string.IsNullOrWhiteSpace(first) || first.Length < 2 || first[0] != 'S' || !char.IsDigit(first[1]))
  273. return null;
  274. return first + "_";
  275. }
  276. private static bool IsBuiltInEntity(MdpEntity entity) =>
  277. entity.EntityCode.StartsWith("S1_", StringComparison.OrdinalIgnoreCase)
  278. || entity.EntityCode.StartsWith("S2_", StringComparison.OrdinalIgnoreCase)
  279. || entity.EntityCode.StartsWith("S3_", StringComparison.OrdinalIgnoreCase)
  280. || entity.EntityCode.StartsWith("S4_", StringComparison.OrdinalIgnoreCase);
  281. private async Task<int> ResolveInboundEnabledAsync(long sourceId, MdpEntityUpsertInput input)
  282. {
  283. if (input.InboundEnabled != 1)
  284. return 0;
  285. var sourceType = await _db.Queryable<MdpSource>()
  286. .Where(s => s.Id == sourceId)
  287. .Select(s => s.SourceType)
  288. .FirstAsync();
  289. if (!string.Equals(sourceType, "API_INBOUND", StringComparison.OrdinalIgnoreCase))
  290. throw Oops.Oh("只有对方推送来源可以开通入站");
  291. if (string.IsNullOrWhiteSpace(input.BizKeyExpr))
  292. throw Oops.Oh("推送实体必须填写业务键");
  293. return 1;
  294. }
  295. private static void ValidateEntityUpsert(MdpEntityUpsertInput input)
  296. {
  297. if (string.IsNullOrWhiteSpace(input.EntityCode))
  298. throw Oops.Oh("实体编码不能为空");
  299. if (string.IsNullOrWhiteSpace(input.EntityName))
  300. throw Oops.Oh("实体名称不能为空");
  301. }
  302. private static void ValidateFieldMappingUpsert(MdpFieldMappingUpsertInput input)
  303. {
  304. if (string.IsNullOrWhiteSpace(input.SourceField))
  305. throw Oops.Oh("源字段不能为空");
  306. if (string.IsNullOrWhiteSpace(input.TargetField))
  307. throw Oops.Oh("目标字段不能为空");
  308. }
  309. private static string NormalizeFieldType(string? fieldType)
  310. {
  311. var normalized = string.IsNullOrWhiteSpace(fieldType) ? "DIRECT" : fieldType.Trim().ToUpperInvariant();
  312. return ValidFieldTypes.Contains(normalized) ? normalized : "DIRECT";
  313. }
  314. private static MdpFieldMapping MapFieldMappingInsert(long entityId, MdpFieldMappingUpsertInput input) =>
  315. new()
  316. {
  317. EntityId = entityId,
  318. SourceField = input.SourceField.Trim(),
  319. TargetField = input.TargetField.Trim(),
  320. FieldType = NormalizeFieldType(input.FieldType),
  321. TransformScript = input.TransformScript?.Trim(),
  322. ConstValue = input.ConstValue?.Trim(),
  323. LookupTable = input.LookupTable?.Trim(),
  324. IsRequired = input.IsRequired ? 1 : 0,
  325. DefaultValue = input.DefaultValue?.Trim(),
  326. SortOrder = input.SortOrder,
  327. CreateTime = DateTime.Now
  328. };
  329. private static MdpFieldMappingRow MapFieldMappingRow(MdpFieldMapping row) =>
  330. new()
  331. {
  332. Id = row.Id,
  333. EntityId = row.EntityId,
  334. SourceField = row.SourceField,
  335. TargetField = row.TargetField,
  336. FieldType = row.FieldType,
  337. TransformScript = row.TransformScript,
  338. ConstValue = row.ConstValue,
  339. LookupTable = row.LookupTable,
  340. IsRequired = row.IsRequired != 0,
  341. DefaultValue = row.DefaultValue,
  342. SortOrder = row.SortOrder
  343. };
  344. }