MdpSyncEntityConfigService.cs 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  2. using Admin.NET.Plugin.AiDOP.Dto.DataPlatform;
  3. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  5. /// <summary>
  6. /// 同步任务实体与字段映射配置 API。
  7. /// </summary>
  8. [ApiDescriptionSettings(Order = 324, Description = "数据中台同步实体与字段映射")]
  9. [Route("api/DataPlatform")]
  10. [NonUnify]
  11. public class MdpSyncEntityConfigService : IDynamicApiController, ITransient
  12. {
  13. private static readonly HashSet<string> ValidFieldTypes = new(StringComparer.OrdinalIgnoreCase)
  14. {
  15. "DIRECT", "JSONPATH", "SCRIPT", "CONST", "LOOKUP"
  16. };
  17. private readonly ISqlSugarClient _db;
  18. private readonly UserManager _userManager;
  19. private readonly MdpSourcePullDispatcher _dispatcher;
  20. public MdpSyncEntityConfigService(
  21. ISqlSugarClient db,
  22. UserManager userManager,
  23. MdpSourcePullDispatcher dispatcher)
  24. {
  25. _db = db;
  26. _userManager = userManager;
  27. _dispatcher = dispatcher;
  28. }
  29. [DisplayName("同步任务实体列表")]
  30. [HttpGet("sync-tasks/{taskCode}/entities")]
  31. public async Task<object> GetEntities(string taskCode)
  32. {
  33. var tenantId = _userManager.TenantId;
  34. var task = await FindTaskAsync(taskCode, tenantId);
  35. var prefix = ResolveEntityCodePrefix(task.TaskCode);
  36. var jobCode = string.IsNullOrWhiteSpace(task.JobCode) ? task.TaskCode : task.JobCode;
  37. var where = new List<string> { "(e.tenant_id = @TenantId OR e.tenant_id = 0)", "e.status = 1" };
  38. var pars = new List<SugarParameter> { new("@TenantId", tenantId) };
  39. if (!string.IsNullOrWhiteSpace(prefix))
  40. {
  41. where.Add("(e.entity_code LIKE @EntityPrefix OR e.job_id = @JobCode)");
  42. pars.Add(new SugarParameter("@EntityPrefix", prefix + "%"));
  43. pars.Add(new SugarParameter("@JobCode", jobCode));
  44. }
  45. else
  46. {
  47. where.Add("e.job_id = @JobCode");
  48. pars.Add(new SugarParameter("@JobCode", jobCode));
  49. }
  50. var whereSql = string.Join(" AND ", where);
  51. var list = await _db.Ado.SqlQueryAsync<MdpEntityRow>(
  52. $"""
  53. SELECT e.id AS Id, e.tenant_id AS TenantId, e.source_id AS SourceId,
  54. s.source_code AS SourceCode, s.source_name AS SourceName, s.source_type AS SourceType,
  55. e.entity_code AS EntityCode, e.entity_name AS EntityName, e.entity_type AS EntityType,
  56. e.source_table_name AS SourceTableName, e.source_api_path AS SourceApiPath,
  57. e.target_table_name AS TargetTableName, e.sync_mode AS SyncMode, e.incr_column AS IncrColumn,
  58. e.batch_size AS BatchSize, e.job_id AS JobId, e.status AS Status, e.remark AS Remark,
  59. IFNULL(m.cnt, 0) AS FieldMappingCount
  60. FROM mdp_entity e
  61. LEFT JOIN mdp_source s ON s.id = e.source_id
  62. LEFT JOIN (
  63. SELECT entity_id, COUNT(1) AS cnt FROM mdp_field_mapping GROUP BY entity_id
  64. ) m ON m.entity_id = e.id
  65. WHERE {whereSql}
  66. ORDER BY e.entity_code
  67. """,
  68. pars);
  69. return list;
  70. }
  71. [DisplayName("新增同步任务实体")]
  72. [HttpPost("sync-tasks/{taskCode}/entities")]
  73. public async Task<object> CreateEntity(string taskCode, [FromBody] MdpEntityUpsertInput input)
  74. {
  75. ValidateEntityUpsert(input);
  76. var tenantId = _userManager.TenantId;
  77. var task = await FindTaskAsync(taskCode, tenantId);
  78. var sourceId = await ResolveSourceIdAsync(input, tenantId);
  79. var entityTenantId = task.TenantId == 0 ? 0L : tenantId;
  80. var jobCode = string.IsNullOrWhiteSpace(task.JobCode) ? task.TaskCode : task.JobCode;
  81. var exists = await _db.Queryable<MdpEntity>()
  82. .Where(u => u.TenantId == entityTenantId && u.EntityCode == input.EntityCode.Trim())
  83. .AnyAsync();
  84. if (exists)
  85. throw Oops.Oh("实体编码已存在");
  86. var now = DateTime.Now;
  87. var entity = new MdpEntity
  88. {
  89. TenantId = entityTenantId,
  90. SourceId = sourceId,
  91. EntityCode = input.EntityCode.Trim(),
  92. EntityName = input.EntityName.Trim(),
  93. EntityType = string.IsNullOrWhiteSpace(input.EntityType) ? "TABLE" : input.EntityType.Trim(),
  94. SourceTableName = input.SourceTableName?.Trim(),
  95. SourceApiPath = input.SourceApiPath?.Trim(),
  96. TargetTableName = input.TargetTableName?.Trim(),
  97. SyncMode = string.IsNullOrWhiteSpace(input.SyncMode) ? "INCR" : input.SyncMode.Trim(),
  98. IncrColumn = input.IncrColumn?.Trim(),
  99. BatchSize = input.BatchSize <= 0 ? 5000 : input.BatchSize,
  100. JobId = string.IsNullOrWhiteSpace(input.JobId) ? jobCode : input.JobId.Trim(),
  101. Status = input.Status <= 0 ? 0 : 1,
  102. Remark = input.Remark?.Trim(),
  103. InboundEnabled = await ResolveInboundEnabledAsync(sourceId, input),
  104. BizKeyExpr = input.BizKeyExpr?.Trim(),
  105. CreateTime = now,
  106. UpdateTime = now
  107. };
  108. var id = await _db.Insertable(entity).ExecuteReturnBigIdentityAsync();
  109. await BumpTaskConfigVersionAsync(task);
  110. return new { id };
  111. }
  112. /// <summary>
  113. /// 按实体立即拉取一次(源 → 贴源)。
  114. ///
  115. /// <para><b>为什么需要它</b>:<see cref="Executors.MdpModuleStagingPuller"/> 的实体清单全部由调用方
  116. /// 显式传入,各模块只传自己硬编码的那一份(<c>S1MdpEntityConfig.All</c> 等),配置层新增的
  117. /// 自定义实体没有任何运行时入口,登记完也拉不动。这里补一个受控的通用触发口。</para>
  118. ///
  119. /// <para><b>越权面收敛</b>:不把 entityCode 直接透传给 dispatcher,而是先按「归属本任务 +
  120. /// 当前租户可见 + 已启用」三条件校验,避免借该端点拉取任意实体。</para>
  121. /// </summary>
  122. [DisplayName("按实体立即拉取一次")]
  123. [HttpPost("sync-tasks/{taskCode}/entities/{entityCode}/pull")]
  124. public async Task<object> PullEntityNow(
  125. string taskCode,
  126. string entityCode,
  127. [FromQuery] bool fullRefresh = false,
  128. CancellationToken cancellationToken = default)
  129. {
  130. var tenantId = _userManager.TenantId;
  131. // tenantId<=0 时贴源写入的租户解析会回落到 null,整行被静默跳过(只打 Warning),
  132. // 对外表现为「拉取成功但 0 行」。提前挡住,不给静默失败留口子。
  133. if (tenantId <= 0)
  134. throw Oops.Oh("无法解析当前租户,请在具体租户下执行拉取");
  135. var task = await FindTaskAsync(taskCode, tenantId);
  136. var jobCode = string.IsNullOrWhiteSpace(task.JobCode) ? task.TaskCode : task.JobCode;
  137. var code = entityCode.Trim();
  138. var target = await _db.Queryable<MdpEntity>()
  139. .Where(u => u.EntityCode == code
  140. && u.Status == 1
  141. && u.JobId == jobCode
  142. && (u.TenantId == tenantId || u.TenantId == 0))
  143. .FirstAsync(cancellationToken)
  144. ?? throw Oops.Oh($"实体 {code} 不存在、未启用,或不归属任务 {taskCode}");
  145. // sync_batch_id 是 varchar(64):固定段 22 字符 + 实体编码最多截 42 字符,
  146. // 截尾只落在实体编码上,不动时间戳,避免批次号互相撞。
  147. var codePart = code.Length > 42 ? code[..42] : code;
  148. var ctx = new MdpPullContext
  149. {
  150. TenantId = tenantId,
  151. FullRefresh = fullRefresh,
  152. BatchId = $"MANUAL_{DateTime.Now:yyyyMMddHHmmss}_{codePart}",
  153. TaskCode = task.TaskCode
  154. };
  155. try
  156. {
  157. var result = await _dispatcher.PullAllByEntityCodeAsync(code, ctx, cancellationToken);
  158. return new
  159. {
  160. entityCode = code,
  161. targetTable = target.TargetTableName,
  162. batchId = ctx.BatchId,
  163. requestedFullRefresh = fullRefresh,
  164. // 调度窗口可能把 FULL 落成 ctx.FullRefresh=true,回传实际生效值,不做黑盒
  165. effectiveWindowType = ctx.SyncWindowType,
  166. effectiveFullRefresh = ctx.FullRefresh,
  167. rowsPulled = result.RowsPulled,
  168. rowsWritten = result.RowsWritten,
  169. newCursor = result.NewCursor,
  170. message = result.Message
  171. };
  172. }
  173. catch (InvalidOperationException ex)
  174. {
  175. // 贴源契约缺列、源表标识符非法、源不可达等都走这里;转业务异常把原因如实带回前端
  176. throw Oops.Oh(ex.Message);
  177. }
  178. }
  179. [DisplayName("更新同步实体")]
  180. [HttpPut("entities/{id:long}")]
  181. public async Task<object> UpdateEntity(long id, [FromBody] MdpEntityUpsertInput input)
  182. {
  183. ValidateEntityUpsert(input);
  184. var tenantId = _userManager.TenantId;
  185. var entity = await FindEditableEntityAsync(id, tenantId);
  186. entity.EntityName = input.EntityName.Trim();
  187. entity.EntityType = string.IsNullOrWhiteSpace(input.EntityType) ? "TABLE" : input.EntityType.Trim();
  188. entity.SourceId = await ResolveSourceIdAsync(input, tenantId, entity.SourceId);
  189. entity.SourceTableName = input.SourceTableName?.Trim();
  190. entity.SourceApiPath = input.SourceApiPath?.Trim();
  191. entity.TargetTableName = input.TargetTableName?.Trim();
  192. entity.SyncMode = string.IsNullOrWhiteSpace(input.SyncMode) ? "INCR" : input.SyncMode.Trim();
  193. entity.IncrColumn = input.IncrColumn?.Trim();
  194. entity.BatchSize = input.BatchSize <= 0 ? 5000 : input.BatchSize;
  195. entity.JobId = input.JobId?.Trim();
  196. entity.Status = input.Status <= 0 ? 0 : 1;
  197. entity.Remark = input.Remark?.Trim();
  198. entity.InboundEnabled = await ResolveInboundEnabledAsync(entity.SourceId, input);
  199. entity.BizKeyExpr = input.BizKeyExpr?.Trim();
  200. entity.UpdateTime = DateTime.Now;
  201. if (!IsBuiltInEntity(entity))
  202. entity.EntityCode = input.EntityCode.Trim();
  203. await _db.Updateable(entity).ExecuteCommandAsync();
  204. return new { id = entity.Id };
  205. }
  206. [DisplayName("删除同步实体")]
  207. [HttpDelete("entities/{id:long}")]
  208. public async Task<object> DeleteEntity(long id)
  209. {
  210. var tenantId = _userManager.TenantId;
  211. var entity = await FindEditableEntityAsync(id, tenantId);
  212. if (IsBuiltInEntity(entity) && entity.TenantId == 0)
  213. throw Oops.Oh("内置 MDP 实体不允许删除");
  214. entity.Status = 0;
  215. entity.UpdateTime = DateTime.Now;
  216. await _db.Updateable(entity).UpdateColumns(u => new { u.Status, u.UpdateTime }).ExecuteCommandAsync();
  217. return new { id = entity.Id };
  218. }
  219. [DisplayName("实体字段映射列表")]
  220. [HttpGet("entities/{entityId:long}/field-mappings")]
  221. public async Task<object> GetFieldMappings(long entityId)
  222. {
  223. var tenantId = _userManager.TenantId;
  224. await FindEditableEntityAsync(entityId, tenantId);
  225. var list = await _db.Queryable<MdpFieldMapping>()
  226. .Where(u => u.EntityId == entityId)
  227. .OrderBy(u => u.SortOrder)
  228. .OrderBy(u => u.TargetField)
  229. .ToListAsync();
  230. return list.Select(MapFieldMappingRow).ToList();
  231. }
  232. [DisplayName("新增实体字段映射")]
  233. [HttpPost("entities/{entityId:long}/field-mappings")]
  234. public async Task<object> CreateFieldMapping(long entityId, [FromBody] MdpFieldMappingUpsertInput input)
  235. {
  236. ValidateFieldMappingUpsert(input);
  237. var tenantId = _userManager.TenantId;
  238. await FindEditableEntityAsync(entityId, tenantId);
  239. var exists = await _db.Queryable<MdpFieldMapping>()
  240. .Where(u => u.EntityId == entityId && u.TargetField == input.TargetField.Trim())
  241. .AnyAsync();
  242. if (exists)
  243. throw Oops.Oh("目标字段映射已存在");
  244. var entity = MapFieldMappingInsert(entityId, input);
  245. var id = await _db.Insertable(entity).ExecuteReturnBigIdentityAsync();
  246. return new { id };
  247. }
  248. [DisplayName("更新字段映射")]
  249. [HttpPut("field-mappings/{id:long}")]
  250. public async Task<object> UpdateFieldMapping(long id, [FromBody] MdpFieldMappingUpsertInput input)
  251. {
  252. ValidateFieldMappingUpsert(input);
  253. var tenantId = _userManager.TenantId;
  254. var mapping = await _db.Queryable<MdpFieldMapping>().Where(u => u.Id == id).FirstAsync();
  255. if (mapping == null)
  256. throw Oops.Oh("字段映射不存在");
  257. await FindEditableEntityAsync(mapping.EntityId, tenantId);
  258. var duplicate = await _db.Queryable<MdpFieldMapping>()
  259. .Where(u => u.EntityId == mapping.EntityId && u.TargetField == input.TargetField.Trim() && u.Id != id)
  260. .AnyAsync();
  261. if (duplicate)
  262. throw Oops.Oh("目标字段映射已存在");
  263. mapping.SourceField = input.SourceField.Trim();
  264. mapping.TargetField = input.TargetField.Trim();
  265. mapping.FieldType = NormalizeFieldType(input.FieldType);
  266. mapping.TransformScript = input.TransformScript?.Trim();
  267. mapping.ConstValue = input.ConstValue?.Trim();
  268. mapping.LookupTable = input.LookupTable?.Trim();
  269. mapping.IsRequired = input.IsRequired ? 1 : 0;
  270. mapping.DefaultValue = input.DefaultValue?.Trim();
  271. mapping.SortOrder = input.SortOrder;
  272. await _db.Updateable(mapping).ExecuteCommandAsync();
  273. return new { id = mapping.Id };
  274. }
  275. [DisplayName("删除字段映射")]
  276. [HttpDelete("field-mappings/{id:long}")]
  277. public async Task<object> DeleteFieldMapping(long id)
  278. {
  279. var tenantId = _userManager.TenantId;
  280. var mapping = await _db.Queryable<MdpFieldMapping>().Where(u => u.Id == id).FirstAsync();
  281. if (mapping == null)
  282. throw Oops.Oh("字段映射不存在");
  283. await FindEditableEntityAsync(mapping.EntityId, tenantId);
  284. await _db.Deleteable<MdpFieldMapping>().Where(u => u.Id == id).ExecuteCommandAsync();
  285. return new { id };
  286. }
  287. private async Task<MdpSyncTask> FindTaskAsync(string taskCode, long tenantId)
  288. {
  289. if (string.IsNullOrWhiteSpace(taskCode))
  290. throw Oops.Oh("任务编码不能为空");
  291. var task = await _db.Queryable<MdpSyncTask>()
  292. .Where(u => u.TaskCode == taskCode && (u.TenantId == tenantId || u.TenantId == 0) && u.Status == 1)
  293. .OrderBy(u => u.TenantId, OrderByType.Desc)
  294. .FirstAsync();
  295. if (task == null)
  296. throw Oops.Oh("同步任务不存在");
  297. return task;
  298. }
  299. private async Task<MdpEntity> FindEditableEntityAsync(long id, long tenantId)
  300. {
  301. var entity = await _db.Queryable<MdpEntity>()
  302. .Where(u => u.Id == id && (u.TenantId == tenantId || u.TenantId == 0))
  303. .FirstAsync();
  304. if (entity == null)
  305. throw Oops.Oh("同步实体不存在");
  306. return entity;
  307. }
  308. private async Task<long> ResolveSourceIdAsync(MdpEntityUpsertInput input, long tenantId, long? fallbackSourceId = null)
  309. {
  310. if (input.SourceId > 0)
  311. return input.SourceId;
  312. if (!string.IsNullOrWhiteSpace(input.SourceCode))
  313. {
  314. var sourceId = await _db.Ado.GetLongAsync(
  315. """
  316. SELECT id FROM mdp_source
  317. WHERE source_code = @SourceCode AND (tenant_id = @TenantId OR tenant_id = 0)
  318. ORDER BY tenant_id DESC
  319. LIMIT 1
  320. """,
  321. new List<SugarParameter>
  322. {
  323. new("@SourceCode", input.SourceCode.Trim()),
  324. new("@TenantId", tenantId)
  325. });
  326. if (sourceId > 0)
  327. return sourceId;
  328. throw Oops.Oh($"数据源 {input.SourceCode} 不存在");
  329. }
  330. if (fallbackSourceId is > 0)
  331. return fallbackSourceId.Value;
  332. throw Oops.Oh("请指定 sourceId 或 sourceCode");
  333. }
  334. private async Task BumpTaskConfigVersionAsync(MdpSyncTask task)
  335. {
  336. task.ConfigVersion += 1;
  337. task.UpdateTime = DateTime.Now;
  338. await _db.Updateable(task).UpdateColumns(u => new { u.ConfigVersion, u.UpdateTime }).ExecuteCommandAsync();
  339. }
  340. private static string? ResolveEntityCodePrefix(string taskCode)
  341. {
  342. if (string.IsNullOrWhiteSpace(taskCode))
  343. return null;
  344. var first = taskCode.Split('_', StringSplitOptions.RemoveEmptyEntries).FirstOrDefault();
  345. if (string.IsNullOrWhiteSpace(first) || first.Length < 2 || first[0] != 'S' || !char.IsDigit(first[1]))
  346. return null;
  347. return first + "_";
  348. }
  349. private static bool IsBuiltInEntity(MdpEntity entity) =>
  350. entity.EntityCode.StartsWith("S1_", StringComparison.OrdinalIgnoreCase)
  351. || entity.EntityCode.StartsWith("S2_", StringComparison.OrdinalIgnoreCase)
  352. || entity.EntityCode.StartsWith("S3_", StringComparison.OrdinalIgnoreCase)
  353. || entity.EntityCode.StartsWith("S4_", StringComparison.OrdinalIgnoreCase);
  354. private async Task<int> ResolveInboundEnabledAsync(long sourceId, MdpEntityUpsertInput input)
  355. {
  356. if (input.InboundEnabled != 1)
  357. return 0;
  358. var sourceType = await _db.Queryable<MdpSource>()
  359. .Where(s => s.Id == sourceId)
  360. .Select(s => s.SourceType)
  361. .FirstAsync();
  362. if (!string.Equals(sourceType, "API_INBOUND", StringComparison.OrdinalIgnoreCase))
  363. throw Oops.Oh("只有对方推送来源可以开通入站");
  364. if (string.IsNullOrWhiteSpace(input.BizKeyExpr))
  365. throw Oops.Oh("推送实体必须填写业务键");
  366. return 1;
  367. }
  368. private static void ValidateEntityUpsert(MdpEntityUpsertInput input)
  369. {
  370. if (string.IsNullOrWhiteSpace(input.EntityCode))
  371. throw Oops.Oh("实体编码不能为空");
  372. if (string.IsNullOrWhiteSpace(input.EntityName))
  373. throw Oops.Oh("实体名称不能为空");
  374. }
  375. private static void ValidateFieldMappingUpsert(MdpFieldMappingUpsertInput input)
  376. {
  377. if (string.IsNullOrWhiteSpace(input.SourceField))
  378. throw Oops.Oh("源字段不能为空");
  379. if (string.IsNullOrWhiteSpace(input.TargetField))
  380. throw Oops.Oh("目标字段不能为空");
  381. }
  382. private static string NormalizeFieldType(string? fieldType)
  383. {
  384. var normalized = string.IsNullOrWhiteSpace(fieldType) ? "DIRECT" : fieldType.Trim().ToUpperInvariant();
  385. return ValidFieldTypes.Contains(normalized) ? normalized : "DIRECT";
  386. }
  387. private static MdpFieldMapping MapFieldMappingInsert(long entityId, MdpFieldMappingUpsertInput input) =>
  388. new()
  389. {
  390. EntityId = entityId,
  391. SourceField = input.SourceField.Trim(),
  392. TargetField = input.TargetField.Trim(),
  393. FieldType = NormalizeFieldType(input.FieldType),
  394. TransformScript = input.TransformScript?.Trim(),
  395. ConstValue = input.ConstValue?.Trim(),
  396. LookupTable = input.LookupTable?.Trim(),
  397. IsRequired = input.IsRequired ? 1 : 0,
  398. DefaultValue = input.DefaultValue?.Trim(),
  399. SortOrder = input.SortOrder,
  400. CreateTime = DateTime.Now
  401. };
  402. private static MdpFieldMappingRow MapFieldMappingRow(MdpFieldMapping row) =>
  403. new()
  404. {
  405. Id = row.Id,
  406. EntityId = row.EntityId,
  407. SourceField = row.SourceField,
  408. TargetField = row.TargetField,
  409. FieldType = row.FieldType,
  410. TransformScript = row.TransformScript,
  411. ConstValue = row.ConstValue,
  412. LookupTable = row.LookupTable,
  413. IsRequired = row.IsRequired != 0,
  414. DefaultValue = row.DefaultValue,
  415. SortOrder = row.SortOrder
  416. };
  417. }