|
|
@@ -1,3 +1,4 @@
|
|
|
+using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
|
|
|
using Admin.NET.Plugin.AiDOP.Dto.DataPlatform;
|
|
|
using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
|
|
|
|
|
|
@@ -18,11 +19,16 @@ public class MdpSyncEntityConfigService : IDynamicApiController, ITransient
|
|
|
|
|
|
private readonly ISqlSugarClient _db;
|
|
|
private readonly UserManager _userManager;
|
|
|
+ private readonly MdpSourcePullDispatcher _dispatcher;
|
|
|
|
|
|
- public MdpSyncEntityConfigService(ISqlSugarClient db, UserManager userManager)
|
|
|
+ public MdpSyncEntityConfigService(
|
|
|
+ ISqlSugarClient db,
|
|
|
+ UserManager userManager,
|
|
|
+ MdpSourcePullDispatcher dispatcher)
|
|
|
{
|
|
|
_db = db;
|
|
|
_userManager = userManager;
|
|
|
+ _dispatcher = dispatcher;
|
|
|
}
|
|
|
|
|
|
[DisplayName("同步任务实体列表")]
|
|
|
@@ -117,6 +123,78 @@ public class MdpSyncEntityConfigService : IDynamicApiController, ITransient
|
|
|
return new { id };
|
|
|
}
|
|
|
|
|
|
+ /// <summary>
|
|
|
+ /// 按实体立即拉取一次(源 → 贴源)。
|
|
|
+ ///
|
|
|
+ /// <para><b>为什么需要它</b>:<see cref="Executors.MdpModuleStagingPuller"/> 的实体清单全部由调用方
|
|
|
+ /// 显式传入,各模块只传自己硬编码的那一份(<c>S1MdpEntityConfig.All</c> 等),配置层新增的
|
|
|
+ /// 自定义实体没有任何运行时入口,登记完也拉不动。这里补一个受控的通用触发口。</para>
|
|
|
+ ///
|
|
|
+ /// <para><b>越权面收敛</b>:不把 entityCode 直接透传给 dispatcher,而是先按「归属本任务 +
|
|
|
+ /// 当前租户可见 + 已启用」三条件校验,避免借该端点拉取任意实体。</para>
|
|
|
+ /// </summary>
|
|
|
+ [DisplayName("按实体立即拉取一次")]
|
|
|
+ [HttpPost("sync-tasks/{taskCode}/entities/{entityCode}/pull")]
|
|
|
+ public async Task<object> PullEntityNow(
|
|
|
+ string taskCode,
|
|
|
+ string entityCode,
|
|
|
+ [FromQuery] bool fullRefresh = false,
|
|
|
+ CancellationToken cancellationToken = default)
|
|
|
+ {
|
|
|
+ var tenantId = _userManager.TenantId;
|
|
|
+ // tenantId<=0 时贴源写入的租户解析会回落到 null,整行被静默跳过(只打 Warning),
|
|
|
+ // 对外表现为「拉取成功但 0 行」。提前挡住,不给静默失败留口子。
|
|
|
+ if (tenantId <= 0)
|
|
|
+ throw Oops.Oh("无法解析当前租户,请在具体租户下执行拉取");
|
|
|
+
|
|
|
+ var task = await FindTaskAsync(taskCode, tenantId);
|
|
|
+ var jobCode = string.IsNullOrWhiteSpace(task.JobCode) ? task.TaskCode : task.JobCode;
|
|
|
+ var code = entityCode.Trim();
|
|
|
+
|
|
|
+ var target = await _db.Queryable<MdpEntity>()
|
|
|
+ .Where(u => u.EntityCode == code
|
|
|
+ && u.Status == 1
|
|
|
+ && u.JobId == jobCode
|
|
|
+ && (u.TenantId == tenantId || u.TenantId == 0))
|
|
|
+ .FirstAsync(cancellationToken)
|
|
|
+ ?? throw Oops.Oh($"实体 {code} 不存在、未启用,或不归属任务 {taskCode}");
|
|
|
+
|
|
|
+ // sync_batch_id 是 varchar(64):固定段 22 字符 + 实体编码最多截 42 字符,
|
|
|
+ // 截尾只落在实体编码上,不动时间戳,避免批次号互相撞。
|
|
|
+ var codePart = code.Length > 42 ? code[..42] : code;
|
|
|
+ var ctx = new MdpPullContext
|
|
|
+ {
|
|
|
+ TenantId = tenantId,
|
|
|
+ FullRefresh = fullRefresh,
|
|
|
+ BatchId = $"MANUAL_{DateTime.Now:yyyyMMddHHmmss}_{codePart}",
|
|
|
+ TaskCode = task.TaskCode
|
|
|
+ };
|
|
|
+
|
|
|
+ try
|
|
|
+ {
|
|
|
+ var result = await _dispatcher.PullAllByEntityCodeAsync(code, ctx, cancellationToken);
|
|
|
+ return new
|
|
|
+ {
|
|
|
+ entityCode = code,
|
|
|
+ targetTable = target.TargetTableName,
|
|
|
+ batchId = ctx.BatchId,
|
|
|
+ requestedFullRefresh = fullRefresh,
|
|
|
+ // 调度窗口可能把 FULL 落成 ctx.FullRefresh=true,回传实际生效值,不做黑盒
|
|
|
+ effectiveWindowType = ctx.SyncWindowType,
|
|
|
+ effectiveFullRefresh = ctx.FullRefresh,
|
|
|
+ rowsPulled = result.RowsPulled,
|
|
|
+ rowsWritten = result.RowsWritten,
|
|
|
+ newCursor = result.NewCursor,
|
|
|
+ message = result.Message
|
|
|
+ };
|
|
|
+ }
|
|
|
+ catch (InvalidOperationException ex)
|
|
|
+ {
|
|
|
+ // 贴源契约缺列、源表标识符非法、源不可达等都走这里;转业务异常把原因如实带回前端
|
|
|
+ throw Oops.Oh(ex.Message);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
[DisplayName("更新同步实体")]
|
|
|
[HttpPut("entities/{id:long}")]
|
|
|
public async Task<object> UpdateEntity(long id, [FromBody] MdpEntityUpsertInput input)
|