MdpSyncTaskConfigService.cs 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680
  1. using Admin.NET.Plugin.AiDOP.Dto.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  3. using Admin.NET.Plugin.AiDOP.Order;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  5. /// <summary>
  6. /// 数据中台同步配置中心 API。
  7. /// </summary>
  8. [ApiDescriptionSettings(Order = 323, Description = "数据中台同步配置中心")]
  9. [Route("api/DataPlatform")]
  10. [AllowAnonymous]
  11. [NonUnify]
  12. public class MdpSyncTaskConfigService : IDynamicApiController, ITransient
  13. {
  14. private static readonly HashSet<string> ProtectedTaskCodes = new(StringComparer.OrdinalIgnoreCase)
  15. {
  16. "S1_MDP_SYNC_TRANSFORM",
  17. "S2_MDP_SYNC_TRANSFORM",
  18. "S3_MDP_SYNC_TRANSFORM",
  19. "S4_MDP_SYNC_TRANSFORM",
  20. "S5_MDP_SYNC_TRANSFORM",
  21. "S5_PURCHASE_RECEIPT_INBOUND",
  22. "S6_MDP_SYNC_TRANSFORM",
  23. "S6_IPQC_INSPECTION_INBOUND",
  24. "S7_MDP_SYNC_TRANSFORM",
  25. "S7_PRODUCTION_RECEIPT_INBOUND"
  26. };
  27. private readonly ISqlSugarClient _db;
  28. private readonly UserManager _userManager;
  29. public MdpSyncTaskConfigService(ISqlSugarClient db, UserManager userManager)
  30. {
  31. _db = db;
  32. _userManager = userManager;
  33. }
  34. [DisplayName("同步任务列表")]
  35. [HttpGet("sync-tasks")]
  36. public async Task<object> GetList([FromQuery] MdpSyncTaskListQuery input)
  37. {
  38. var tenantId = _userManager.TenantId;
  39. var page = input.Page <= 0 ? 1 : input.Page;
  40. var pageSize = input.PageSize <= 0 ? 10 : input.PageSize;
  41. var offset = (page - 1) * pageSize;
  42. var where = new List<string> { "(t.tenant_id = @TenantId OR t.tenant_id = 0)", "t.status = 1" };
  43. var pars = new List<SugarParameter> { new("@TenantId", tenantId) };
  44. if (!string.IsNullOrWhiteSpace(input.Keyword))
  45. {
  46. where.Add("(t.task_code LIKE @Keyword OR t.task_name LIKE @Keyword OR t.business_domain_name LIKE @Keyword)");
  47. pars.Add(new SugarParameter("@Keyword", $"%{input.Keyword.Trim()}%"));
  48. }
  49. if (!string.IsNullOrWhiteSpace(input.BusinessDomainCode))
  50. {
  51. where.Add("t.business_domain_code = @BusinessDomainCode");
  52. pars.Add(new SugarParameter("@BusinessDomainCode", input.BusinessDomainCode.Trim()));
  53. }
  54. if (input.Status.HasValue)
  55. {
  56. where[1] = "t.status = @Status";
  57. pars.Add(new SugarParameter("@Status", input.Status.Value));
  58. }
  59. var whereSql = string.Join(" AND ", where);
  60. var total = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM mdp_sync_task t WHERE {whereSql}", pars);
  61. var list = await _db.Ado.SqlQueryAsync<MdpSyncTaskListRow>(
  62. $"""
  63. SELECT t.id AS Id, t.tenant_id AS TenantId, t.task_code AS TaskCode, t.task_name AS TaskName,
  64. t.task_type AS TaskType, t.business_domain_code AS BusinessDomainCode,
  65. t.business_domain_name AS BusinessDomainName, t.consumer_modules AS ConsumerModules,
  66. t.source_system_code AS SourceSystemCode, t.service_key AS ServiceKey,
  67. t.job_code AS JobCode, t.schedule_job_id AS ScheduleJobId, t.status AS Status,
  68. t.config_version AS ConfigVersion, t.description AS Description,
  69. sch.schedule_mode AS ScheduleMode,
  70. IFNULL(sch.auto_enabled, 1) AS AutoEnabled,
  71. IFNULL(sch.manual_enabled, 1) AS ManualEnabled,
  72. run.status AS LastRunStatus, run.start_time AS LastRunTime
  73. FROM mdp_sync_task t
  74. LEFT JOIN mdp_sync_task_schedule sch
  75. ON sch.tenant_id = t.tenant_id AND sch.task_code = t.task_code
  76. LEFT JOIN (
  77. SELECT job_code, MAX(start_time) AS last_run_time
  78. FROM mdp_transform_run_log
  79. WHERE {MdpMonitorService.BuildMdpRunLogTenantWhere(tenantId)}
  80. GROUP BY job_code
  81. ) latest ON latest.job_code = IFNULL(t.job_code, t.task_code)
  82. LEFT JOIN mdp_transform_run_log run
  83. ON run.job_code = latest.job_code AND run.start_time = latest.last_run_time
  84. WHERE {whereSql}
  85. ORDER BY t.task_code
  86. LIMIT {pageSize} OFFSET {offset}
  87. """,
  88. pars);
  89. return new { total, page, pageSize, list };
  90. }
  91. [DisplayName("同步任务详情")]
  92. [HttpGet("sync-tasks/{id:long}")]
  93. public async Task<object> GetDetail(long id)
  94. {
  95. var tenantId = _userManager.TenantId;
  96. var row = await _db.Ado.SqlQuerySingleAsync<MdpSyncTaskListRow>(
  97. $"""
  98. SELECT t.id AS Id, t.tenant_id AS TenantId, t.task_code AS TaskCode, t.task_name AS TaskName,
  99. t.task_type AS TaskType, t.business_domain_code AS BusinessDomainCode,
  100. t.business_domain_name AS BusinessDomainName, t.consumer_modules AS ConsumerModules,
  101. t.source_system_code AS SourceSystemCode, t.service_key AS ServiceKey,
  102. t.job_code AS JobCode, t.schedule_job_id AS ScheduleJobId, t.status AS Status,
  103. t.config_version AS ConfigVersion, t.description AS Description,
  104. sch.schedule_mode AS ScheduleMode,
  105. IFNULL(sch.auto_enabled, 1) AS AutoEnabled,
  106. IFNULL(sch.manual_enabled, 1) AS ManualEnabled,
  107. run.status AS LastRunStatus, run.start_time AS LastRunTime
  108. FROM mdp_sync_task t
  109. LEFT JOIN mdp_sync_task_schedule sch
  110. ON sch.tenant_id = t.tenant_id AND sch.task_code = t.task_code
  111. LEFT JOIN (
  112. SELECT job_code, MAX(start_time) AS last_run_time
  113. FROM mdp_transform_run_log
  114. WHERE {MdpMonitorService.BuildMdpRunLogTenantWhere(tenantId)}
  115. GROUP BY job_code
  116. ) latest ON latest.job_code = IFNULL(t.job_code, t.task_code)
  117. LEFT JOIN mdp_transform_run_log run
  118. ON run.job_code = latest.job_code AND run.start_time = latest.last_run_time
  119. WHERE t.id = @Id AND (t.tenant_id = @TenantId OR t.tenant_id = 0)
  120. LIMIT 1
  121. """,
  122. new SugarParameter("@Id", id),
  123. new SugarParameter("@TenantId", tenantId));
  124. if (row == null)
  125. throw Oops.Oh("同步任务不存在");
  126. return row;
  127. }
  128. [DisplayName("新增同步任务")]
  129. [HttpPost("sync-tasks")]
  130. public async Task<object> Create([FromBody] MdpSyncTaskUpsertInput input)
  131. {
  132. ValidateUpsert(input);
  133. var tenantId = _userManager.TenantId;
  134. var exists = await _db.Queryable<MdpSyncTask>()
  135. .Where(u => u.TenantId == tenantId && u.TaskCode == input.TaskCode.Trim())
  136. .AnyAsync();
  137. if (exists)
  138. throw Oops.Oh("任务编码已存在");
  139. var entity = MapUpsert(input, tenantId);
  140. entity.ConfigVersion = 1;
  141. entity.CreateTime = DateTime.Now;
  142. entity.UpdateTime = DateTime.Now;
  143. var id = await _db.Insertable(entity).ExecuteReturnBigIdentityAsync();
  144. await EnsureDefaultScheduleAsync(tenantId, entity.TaskCode, entity.ScheduleJobId);
  145. return new { id };
  146. }
  147. [DisplayName("更新同步任务")]
  148. [HttpPut("sync-tasks/{id:long}")]
  149. public async Task<object> Update(long id, [FromBody] MdpSyncTaskUpsertInput input)
  150. {
  151. ValidateUpsert(input);
  152. var tenantId = _userManager.TenantId;
  153. var entity = await FindEditableTaskAsync(id, tenantId);
  154. entity.TaskName = input.TaskName.Trim();
  155. entity.TaskType = string.IsNullOrWhiteSpace(input.TaskType) ? "SERVICE_SYNC" : input.TaskType.Trim();
  156. entity.BusinessDomainCode = input.BusinessDomainCode?.Trim();
  157. entity.BusinessDomainName = input.BusinessDomainName?.Trim();
  158. entity.ConsumerModules = input.ConsumerModules?.Trim();
  159. entity.SourceSystemCode = input.SourceSystemCode?.Trim();
  160. entity.ServiceKey = input.ServiceKey?.Trim();
  161. entity.JobCode = string.IsNullOrWhiteSpace(input.JobCode) ? input.TaskCode.Trim() : input.JobCode.Trim();
  162. entity.ScheduleJobId = input.ScheduleJobId?.Trim();
  163. entity.Status = input.Status <= 0 ? 0 : 1;
  164. entity.OwnerRole = input.OwnerRole?.Trim();
  165. entity.Description = input.Description?.Trim();
  166. entity.ConfigVersion += 1;
  167. entity.UpdateTime = DateTime.Now;
  168. if (!ProtectedTaskCodes.Contains(entity.TaskCode))
  169. entity.TaskCode = input.TaskCode.Trim();
  170. await _db.Updateable(entity).ExecuteCommandAsync();
  171. await SyncScheduleJobIdAsync(entity.TenantId, entity.TaskCode, entity.ScheduleJobId);
  172. return new { id = entity.Id };
  173. }
  174. [DisplayName("删除同步任务")]
  175. [HttpDelete("sync-tasks/{id:long}")]
  176. public async Task<object> Delete(long id)
  177. {
  178. var tenantId = _userManager.TenantId;
  179. var entity = await FindEditableTaskAsync(id, tenantId);
  180. if (ProtectedTaskCodes.Contains(entity.TaskCode) && entity.TenantId == 0)
  181. throw Oops.Oh("内置 MDP 同步任务不允许删除");
  182. entity.Status = 0;
  183. entity.UpdateTime = DateTime.Now;
  184. await _db.Updateable(entity).UpdateColumns(u => new { u.Status, u.UpdateTime }).ExecuteCommandAsync();
  185. return new { id };
  186. }
  187. [DisplayName("同步任务步骤列表")]
  188. [HttpGet("sync-tasks/{taskCode}/steps")]
  189. public async Task<object> GetSteps(string taskCode)
  190. {
  191. var tenantId = _userManager.TenantId;
  192. await EnsureTaskExistsAsync(taskCode, tenantId);
  193. var steps = await _db.Queryable<MdpSyncTaskStep>()
  194. .Where(u => u.TaskCode == taskCode && (u.TenantId == tenantId || u.TenantId == 0))
  195. .OrderBy(u => u.SortOrder)
  196. .ToListAsync();
  197. return steps.Select(MapStepRow).ToList();
  198. }
  199. [DisplayName("保存同步任务步骤")]
  200. [HttpPut("sync-tasks/{taskCode}/steps")]
  201. public async Task<object> SaveSteps(string taskCode, [FromBody] List<MdpSyncTaskStepRow> input)
  202. {
  203. var tenantId = _userManager.TenantId;
  204. var task = await EnsureTaskExistsAsync(taskCode, tenantId);
  205. var rows = input ?? new List<MdpSyncTaskStepRow>();
  206. var now = DateTime.Now;
  207. foreach (var row in rows)
  208. {
  209. if (string.IsNullOrWhiteSpace(row.StepCode) || string.IsNullOrWhiteSpace(row.StepName))
  210. throw Oops.Oh("步骤编码和名称不能为空");
  211. var stage = string.IsNullOrWhiteSpace(row.StageType) ? row.StepCode.Trim() : row.StageType.Trim();
  212. if (row.Enabled && Admin.NET.Plugin.AiDOP.DataPlatform.FileImport.MdpExcelValidation.IsPullStage(stage, row.StepCode))
  213. {
  214. if (!string.IsNullOrWhiteSpace(task.SourceSystemCode))
  215. {
  216. var src = await _db.Queryable<MdpSource>().FirstAsync(x => x.SourceCode == task.SourceSystemCode);
  217. if (src != null && string.Equals(src.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase))
  218. throw Oops.Oh("FILE_EXCEL 实体不得登记 Pull/Extract 步骤,请使用文件导入接口");
  219. }
  220. }
  221. var entity = await _db.Queryable<MdpSyncTaskStep>()
  222. .Where(u => u.TaskCode == taskCode && u.StepCode == row.StepCode && (u.TenantId == tenantId || u.TenantId == 0))
  223. .FirstAsync();
  224. if (entity == null)
  225. {
  226. entity = new MdpSyncTaskStep
  227. {
  228. TenantId = task.TenantId == 0 ? 0 : tenantId,
  229. TaskCode = taskCode,
  230. StepCode = row.StepCode.Trim(),
  231. StepName = row.StepName.Trim(),
  232. StageType = string.IsNullOrWhiteSpace(row.StageType) ? row.StepCode.Trim() : row.StageType.Trim(),
  233. ServiceMethodKey = row.ServiceMethodKey?.Trim(),
  234. Enabled = row.Enabled ? 1 : 0,
  235. SortOrder = row.SortOrder,
  236. Description = row.Description?.Trim(),
  237. CreateTime = now,
  238. UpdateTime = now
  239. };
  240. await _db.Insertable(entity).ExecuteCommandAsync();
  241. }
  242. else
  243. {
  244. entity.StepName = row.StepName.Trim();
  245. entity.StageType = string.IsNullOrWhiteSpace(row.StageType) ? row.StepCode.Trim() : row.StageType.Trim();
  246. entity.ServiceMethodKey = row.ServiceMethodKey?.Trim();
  247. entity.Enabled = row.Enabled ? 1 : 0;
  248. entity.SortOrder = row.SortOrder;
  249. entity.Description = row.Description?.Trim();
  250. entity.UpdateTime = now;
  251. await _db.Updateable(entity).ExecuteCommandAsync();
  252. }
  253. }
  254. task.ConfigVersion += 1;
  255. task.UpdateTime = now;
  256. await _db.Updateable(task).UpdateColumns(u => new { u.ConfigVersion, u.UpdateTime }).ExecuteCommandAsync();
  257. return new { count = rows.Count };
  258. }
  259. [DisplayName("同步任务调度配置")]
  260. [HttpGet("sync-tasks/{taskCode}/schedule")]
  261. public async Task<object> GetSchedule(string taskCode)
  262. {
  263. var tenantId = _userManager.TenantId;
  264. await EnsureTaskExistsAsync(taskCode, tenantId);
  265. var schedule = await FindScheduleAsync(taskCode, tenantId);
  266. if (schedule == null)
  267. return new MdpSyncTaskScheduleRow { TaskCode = taskCode, ScheduleMode = "CRON" };
  268. return MapScheduleRow(schedule);
  269. }
  270. [DisplayName("保存同步任务调度配置")]
  271. [HttpPut("sync-tasks/{taskCode}/schedule")]
  272. public async Task<object> SaveSchedule(string taskCode, [FromBody] MdpSyncTaskScheduleRow input)
  273. {
  274. var tenantId = _userManager.TenantId;
  275. var task = await EnsureTaskExistsAsync(taskCode, tenantId);
  276. var schedule = await FindScheduleAsync(taskCode, tenantId);
  277. var now = DateTime.Now;
  278. if (schedule == null)
  279. {
  280. schedule = new MdpSyncTaskSchedule
  281. {
  282. TenantId = task.TenantId == 0 ? 0 : tenantId,
  283. TaskCode = taskCode,
  284. CreateTime = now,
  285. UpdateTime = now
  286. };
  287. }
  288. schedule.ScheduleJobId = input.ScheduleJobId?.Trim() ?? task.ScheduleJobId;
  289. schedule.ScheduleMode = string.IsNullOrWhiteSpace(input.ScheduleMode) ? "CRON" : input.ScheduleMode.Trim();
  290. schedule.CronExpr = input.CronExpr?.Trim();
  291. schedule.CronDesc = input.CronDesc?.Trim();
  292. schedule.Timezone = string.IsNullOrWhiteSpace(input.Timezone) ? "Asia/Shanghai" : input.Timezone.Trim();
  293. schedule.AutoEnabled = input.AutoEnabled ? 1 : 0;
  294. schedule.ManualEnabled = input.ManualEnabled ? 1 : 0;
  295. schedule.RetryEnabled = input.RetryEnabled ? 1 : 0;
  296. schedule.MaxRetryCount = input.MaxRetryCount <= 0 ? 3 : input.MaxRetryCount;
  297. schedule.RetryIntervalSeconds = input.RetryIntervalSeconds <= 0 ? 300 : input.RetryIntervalSeconds;
  298. schedule.TimeoutSeconds = input.TimeoutSeconds <= 0 ? 3600 : input.TimeoutSeconds;
  299. schedule.MisfirePolicy = input.MisfirePolicy?.Trim();
  300. schedule.SyncWindowType = string.IsNullOrWhiteSpace(input.SyncWindowType) ? "FULL" : input.SyncWindowType.Trim();
  301. schedule.SyncWindowValue = input.SyncWindowValue?.Trim();
  302. schedule.LastScheduleTime = input.LastScheduleTime;
  303. schedule.NextScheduleTime = input.NextScheduleTime;
  304. schedule.AdminJobConfigJson = input.AdminJobConfigJson;
  305. schedule.Description = input.Description?.Trim();
  306. schedule.UpdateTime = now;
  307. if (schedule.Id <= 0)
  308. await _db.Insertable(schedule).ExecuteCommandAsync();
  309. else
  310. await _db.Updateable(schedule).ExecuteCommandAsync();
  311. task.ScheduleJobId = schedule.ScheduleJobId;
  312. task.ConfigVersion += 1;
  313. task.UpdateTime = now;
  314. await _db.Updateable(task).UpdateColumns(u => new { u.ScheduleJobId, u.ConfigVersion, u.UpdateTime }).ExecuteCommandAsync();
  315. return MapScheduleRow(schedule);
  316. }
  317. [DisplayName("同步 Admin.NET 调度快照")]
  318. [HttpPost("sync-tasks/{taskCode}/schedule/sync-admin-job")]
  319. public async Task<object> SyncAdminJob(string taskCode)
  320. {
  321. var tenantId = _userManager.TenantId;
  322. var task = await EnsureTaskExistsAsync(taskCode, tenantId);
  323. var schedule = await FindScheduleAsync(taskCode, tenantId) ?? new MdpSyncTaskSchedule
  324. {
  325. TenantId = task.TenantId == 0 ? 0 : tenantId,
  326. TaskCode = taskCode,
  327. ScheduleJobId = task.ScheduleJobId,
  328. ScheduleMode = "CRON",
  329. CreateTime = DateTime.Now,
  330. UpdateTime = DateTime.Now
  331. };
  332. schedule.AdminJobConfigJson = System.Text.Json.JsonSerializer.Serialize(new
  333. {
  334. scheduleJobId = schedule.ScheduleJobId ?? task.ScheduleJobId,
  335. message = "第一阶段仅保存 Admin.NET Job 绑定快照,不写回底层 Cron 配置。"
  336. });
  337. schedule.UpdateTime = DateTime.Now;
  338. if (schedule.Id <= 0)
  339. await _db.Insertable(schedule).ExecuteCommandAsync();
  340. else
  341. await _db.Updateable(schedule).UpdateColumns(u => new { u.AdminJobConfigJson, u.UpdateTime }).ExecuteCommandAsync();
  342. return new
  343. {
  344. taskCode,
  345. scheduleJobId = schedule.ScheduleJobId ?? task.ScheduleJobId,
  346. adminJobConfigJson = schedule.AdminJobConfigJson
  347. };
  348. }
  349. private async Task<MdpSyncTask> EnsureTaskExistsAsync(string taskCode, long tenantId)
  350. {
  351. if (string.IsNullOrWhiteSpace(taskCode))
  352. throw Oops.Oh("任务编码不能为空");
  353. var task = await _db.Queryable<MdpSyncTask>()
  354. .Where(u => u.TaskCode == taskCode && (u.TenantId == tenantId || u.TenantId == 0) && u.Status == 1)
  355. .OrderBy(u => u.TenantId, OrderByType.Desc)
  356. .FirstAsync();
  357. if (task == null)
  358. throw Oops.Oh("同步任务不存在");
  359. return task;
  360. }
  361. private async Task<MdpSyncTask> FindEditableTaskAsync(long id, long tenantId)
  362. {
  363. var entity = await _db.Queryable<MdpSyncTask>()
  364. .Where(u => u.Id == id && (u.TenantId == tenantId || u.TenantId == 0))
  365. .FirstAsync();
  366. if (entity == null)
  367. throw Oops.Oh("同步任务不存在");
  368. return entity;
  369. }
  370. private async Task<MdpSyncTaskSchedule?> FindScheduleAsync(string taskCode, long tenantId) =>
  371. await _db.Queryable<MdpSyncTaskSchedule>()
  372. .Where(u => u.TaskCode == taskCode && (u.TenantId == tenantId || u.TenantId == 0))
  373. .OrderBy(u => u.TenantId, OrderByType.Desc)
  374. .FirstAsync();
  375. private async Task EnsureDefaultScheduleAsync(long tenantId, string taskCode, string? scheduleJobId)
  376. {
  377. var exists = await _db.Queryable<MdpSyncTaskSchedule>()
  378. .Where(u => u.TaskCode == taskCode && (u.TenantId == tenantId || u.TenantId == 0))
  379. .AnyAsync();
  380. if (exists)
  381. return;
  382. await _db.Insertable(new MdpSyncTaskSchedule
  383. {
  384. TenantId = tenantId,
  385. TaskCode = taskCode,
  386. ScheduleJobId = scheduleJobId,
  387. ScheduleMode = "CRON",
  388. CronDesc = "按 Admin.NET 任务调度执行",
  389. AutoEnabled = 1,
  390. ManualEnabled = 1,
  391. RetryEnabled = 1,
  392. MaxRetryCount = 3,
  393. RetryIntervalSeconds = 300,
  394. TimeoutSeconds = 3600,
  395. SyncWindowType = "FULL",
  396. CreateTime = DateTime.Now,
  397. UpdateTime = DateTime.Now
  398. }).ExecuteCommandAsync();
  399. }
  400. private async Task SyncScheduleJobIdAsync(long tenantId, string taskCode, string? scheduleJobId)
  401. {
  402. var schedule = await FindScheduleAsync(taskCode, tenantId);
  403. if (schedule == null)
  404. {
  405. await EnsureDefaultScheduleAsync(tenantId, taskCode, scheduleJobId);
  406. return;
  407. }
  408. schedule.ScheduleJobId = scheduleJobId;
  409. schedule.UpdateTime = DateTime.Now;
  410. await _db.Updateable(schedule).UpdateColumns(u => new { u.ScheduleJobId, u.UpdateTime }).ExecuteCommandAsync();
  411. }
  412. private static void ValidateUpsert(MdpSyncTaskUpsertInput input)
  413. {
  414. if (string.IsNullOrWhiteSpace(input.TaskCode))
  415. throw Oops.Oh("任务编码不能为空");
  416. if (string.IsNullOrWhiteSpace(input.TaskName))
  417. throw Oops.Oh("任务名称不能为空");
  418. }
  419. private static MdpSyncTask MapUpsert(MdpSyncTaskUpsertInput input, long tenantId) =>
  420. new()
  421. {
  422. TenantId = tenantId,
  423. TaskCode = input.TaskCode.Trim(),
  424. TaskName = input.TaskName.Trim(),
  425. TaskType = string.IsNullOrWhiteSpace(input.TaskType) ? "SERVICE_SYNC" : input.TaskType.Trim(),
  426. BusinessDomainCode = input.BusinessDomainCode?.Trim(),
  427. BusinessDomainName = input.BusinessDomainName?.Trim(),
  428. ConsumerModules = input.ConsumerModules?.Trim(),
  429. SourceSystemCode = input.SourceSystemCode?.Trim(),
  430. ServiceKey = input.ServiceKey?.Trim(),
  431. JobCode = string.IsNullOrWhiteSpace(input.JobCode) ? input.TaskCode.Trim() : input.JobCode.Trim(),
  432. ScheduleJobId = input.ScheduleJobId?.Trim(),
  433. Status = input.Status <= 0 ? 0 : 1,
  434. OwnerRole = input.OwnerRole?.Trim(),
  435. Description = input.Description?.Trim()
  436. };
  437. private static MdpSyncTaskStepRow MapStepRow(MdpSyncTaskStep step) =>
  438. new()
  439. {
  440. Id = step.Id,
  441. StepCode = step.StepCode,
  442. StepName = step.StepName,
  443. StageType = step.StageType,
  444. ServiceMethodKey = step.ServiceMethodKey,
  445. Enabled = step.Enabled != 0,
  446. SortOrder = step.SortOrder,
  447. Description = step.Description
  448. };
  449. private static MdpSyncTaskScheduleRow MapScheduleRow(MdpSyncTaskSchedule schedule) =>
  450. new()
  451. {
  452. TaskCode = schedule.TaskCode,
  453. ScheduleJobId = schedule.ScheduleJobId,
  454. ScheduleMode = schedule.ScheduleMode,
  455. CronExpr = schedule.CronExpr,
  456. CronDesc = schedule.CronDesc,
  457. Timezone = schedule.Timezone,
  458. AutoEnabled = schedule.AutoEnabled != 0,
  459. ManualEnabled = schedule.ManualEnabled != 0,
  460. RetryEnabled = schedule.RetryEnabled != 0,
  461. MaxRetryCount = schedule.MaxRetryCount,
  462. RetryIntervalSeconds = schedule.RetryIntervalSeconds,
  463. TimeoutSeconds = schedule.TimeoutSeconds,
  464. MisfirePolicy = schedule.MisfirePolicy,
  465. SyncWindowType = schedule.SyncWindowType,
  466. SyncWindowValue = schedule.SyncWindowValue,
  467. LastScheduleTime = schedule.LastScheduleTime,
  468. NextScheduleTime = schedule.NextScheduleTime,
  469. AdminJobConfigJson = schedule.AdminJobConfigJson,
  470. Description = schedule.Description
  471. };
  472. [DisplayName("任务参数列表")]
  473. [HttpGet("sync-tasks/{taskCode}/params")]
  474. public async Task<object> GetParams(string taskCode)
  475. {
  476. var tenantId = _userManager.TenantId;
  477. var list = await _db.Queryable<MdpSyncTaskParam>()
  478. .Where(u => u.TaskCode == taskCode && (u.TenantId == tenantId || u.TenantId == 0))
  479. .OrderBy(u => u.SortOrder)
  480. .ToListAsync();
  481. return new
  482. {
  483. list = list.Select(p => new
  484. {
  485. p.Id,
  486. p.TaskCode,
  487. p.ScopeType,
  488. p.ScopeCode,
  489. p.ParamKey,
  490. p.ParamName,
  491. p.ParamType,
  492. p.ParamValue,
  493. p.DefaultValue,
  494. required = p.Required != 0,
  495. editable = p.Editable != 0,
  496. p.SortOrder,
  497. p.Description
  498. })
  499. };
  500. }
  501. [DisplayName("保存任务参数")]
  502. [HttpPut("sync-tasks/{taskCode}/params")]
  503. public async Task<object> SaveParams(string taskCode, [FromBody] List<MdpSyncTaskParamSaveItem> items)
  504. {
  505. var tenantId = _userManager.TenantId;
  506. items ??= [];
  507. await _db.Deleteable<MdpSyncTaskParam>()
  508. .Where(u => u.TaskCode == taskCode && u.TenantId == tenantId)
  509. .ExecuteCommandAsync();
  510. var now = DateTime.Now;
  511. var rows = items.Where(x => !string.IsNullOrWhiteSpace(x.ParamKey)).Select((x, i) => new MdpSyncTaskParam
  512. {
  513. TenantId = tenantId,
  514. TaskCode = taskCode,
  515. ScopeType = string.IsNullOrWhiteSpace(x.ScopeType) ? "TASK" : x.ScopeType.Trim(),
  516. ScopeCode = x.ScopeCode?.Trim() ?? "",
  517. ParamKey = x.ParamKey.Trim(),
  518. ParamName = string.IsNullOrWhiteSpace(x.ParamName) ? x.ParamKey.Trim() : x.ParamName.Trim(),
  519. ParamType = string.IsNullOrWhiteSpace(x.ParamType) ? "STRING" : x.ParamType.Trim(),
  520. ParamValue = x.ParamValue,
  521. DefaultValue = x.DefaultValue,
  522. Required = x.Required ? (byte)1 : (byte)0,
  523. Editable = x.Editable ? (byte)1 : (byte)0,
  524. SortOrder = x.SortOrder > 0 ? x.SortOrder : i + 1,
  525. Description = x.Description,
  526. CreateTime = now,
  527. UpdateTime = now
  528. }).ToList();
  529. if (rows.Count > 0)
  530. await _db.Insertable(rows).ExecuteCommandAsync();
  531. return new { ok = true, count = rows.Count };
  532. }
  533. [DisplayName("任务公式列表")]
  534. [HttpGet("sync-tasks/{taskCode}/formulas")]
  535. public async Task<object> GetFormulas(string taskCode)
  536. {
  537. var tenantId = _userManager.TenantId;
  538. var list = await _db.Queryable<MdpSyncTaskFormula>()
  539. .Where(u => u.TaskCode == taskCode && (u.TenantId == tenantId || u.TenantId == 0))
  540. .OrderBy(u => u.SortOrder)
  541. .ToListAsync();
  542. return new
  543. {
  544. list = list.Select(f => new
  545. {
  546. f.Id,
  547. f.TaskCode,
  548. f.StepCode,
  549. f.FormulaCode,
  550. f.FormulaName,
  551. f.MetricCode,
  552. f.FormulaExpr,
  553. f.FormulaPreview,
  554. f.Direction,
  555. f.YellowThreshold,
  556. f.RedThreshold,
  557. f.VersionNo,
  558. enabled = f.IsEnabled != 0,
  559. f.SortOrder,
  560. f.Description
  561. })
  562. };
  563. }
  564. [DisplayName("保存任务公式")]
  565. [HttpPut("sync-tasks/{taskCode}/formulas")]
  566. public async Task<object> SaveFormulas(string taskCode, [FromBody] List<MdpSyncTaskFormulaSaveItem> items)
  567. {
  568. var tenantId = _userManager.TenantId;
  569. items ??= [];
  570. await _db.Deleteable<MdpSyncTaskFormula>()
  571. .Where(u => u.TaskCode == taskCode && u.TenantId == tenantId)
  572. .ExecuteCommandAsync();
  573. var now = DateTime.Now;
  574. var rows = items.Where(x => !string.IsNullOrWhiteSpace(x.FormulaCode)).Select((x, i) => new MdpSyncTaskFormula
  575. {
  576. TenantId = tenantId,
  577. TaskCode = taskCode,
  578. StepCode = x.StepCode,
  579. FormulaCode = x.FormulaCode.Trim(),
  580. FormulaName = string.IsNullOrWhiteSpace(x.FormulaName) ? x.FormulaCode.Trim() : x.FormulaName.Trim(),
  581. MetricCode = x.MetricCode,
  582. FormulaExpr = x.FormulaExpr,
  583. FormulaPreview = x.FormulaPreview,
  584. Direction = string.IsNullOrWhiteSpace(x.Direction) ? "higher_is_better" : x.Direction.Trim(),
  585. YellowThreshold = x.YellowThreshold,
  586. RedThreshold = x.RedThreshold,
  587. VersionNo = x.VersionNo > 0 ? x.VersionNo : 1,
  588. IsEnabled = x.Enabled ? (byte)1 : (byte)0,
  589. SortOrder = x.SortOrder > 0 ? x.SortOrder : i + 1,
  590. Description = x.Description,
  591. CreateTime = now,
  592. UpdateTime = now
  593. }).ToList();
  594. if (rows.Count > 0)
  595. await _db.Insertable(rows).ExecuteCommandAsync();
  596. return new { ok = true, count = rows.Count };
  597. }
  598. }
  599. public class MdpSyncTaskParamSaveItem
  600. {
  601. public string? ScopeType { get; set; }
  602. public string? ScopeCode { get; set; }
  603. public string ParamKey { get; set; } = string.Empty;
  604. public string? ParamName { get; set; }
  605. public string? ParamType { get; set; }
  606. public string? ParamValue { get; set; }
  607. public string? DefaultValue { get; set; }
  608. public bool Required { get; set; }
  609. public bool Editable { get; set; } = true;
  610. public int SortOrder { get; set; }
  611. public string? Description { get; set; }
  612. }
  613. public class MdpSyncTaskFormulaSaveItem
  614. {
  615. public string? StepCode { get; set; }
  616. public string FormulaCode { get; set; } = string.Empty;
  617. public string? FormulaName { get; set; }
  618. public string? MetricCode { get; set; }
  619. public string? FormulaExpr { get; set; }
  620. public string? FormulaPreview { get; set; }
  621. public string? Direction { get; set; }
  622. public decimal? YellowThreshold { get; set; }
  623. public decimal? RedThreshold { get; set; }
  624. public int VersionNo { get; set; } = 1;
  625. public bool Enabled { get; set; } = true;
  626. public int SortOrder { get; set; }
  627. public string? Description { get; set; }
  628. }