using Admin.NET.Plugin.AiDOP.DataPlatform.S1Refresh; using Admin.NET.Plugin.AiDOP.Order; using Furion.Schedule; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using System.Text.Json; namespace Admin.NET.Plugin.AiDOP.Job; /// /// S1 首批 MDP 同步与标准化转换定时任务:按租户/工厂作用域分别入队,由 worker 受控并行。 /// [JobDetail("job_s1_mdp_sync_transform", Description = "S1 MDP同步与标准化转换(含订单交付域原子聚合层)", GroupName = "default", Concurrent = false)] [Period(3600000, TriggerId = "trigger_s1_mdp_sync_transform", Description = "每60分钟执行")] public class S1MdpSyncTransformJob : IJob { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public S1MdpSyncTransformJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory) { _scopeFactory = scopeFactory; _logger = loggerFactory.CreateLogger(nameof(S1MdpSyncTransformJob)); } public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken) { using var scope = _scopeFactory.CreateScope(); var catalog = scope.ServiceProvider.GetRequiredService(); var rebuild = scope.ServiceProvider.GetRequiredService(); var scopes = await catalog.ListEnabledScopesAsync(stoppingToken); var accepted = 0; var skipped = 0; foreach (var runScope in scopes) { stoppingToken.ThrowIfCancellationRequested(); try { var (status, body) = await rebuild.EnqueueAsync(runScope.TenantId, runScope.FactoryId, null, "AUTO", stoppingToken); if (status == 202) accepted++; else skipped++; _logger.LogInformation( "S1MdpSyncTransformJob 作用域 tenant={TenantId} factory={FactoryId} status={Status} jobId={JobId}", runScope.TenantId, runScope.FactoryId, status, body.JobId); } catch (Exception ex) { _logger.LogError(ex, "S1MdpSyncTransformJob 作用域失败 tenant={TenantId} factory={FactoryId}", runScope.TenantId, runScope.FactoryId); } } _logger.LogInformation("S1MdpSyncTransformJob 完成 accepted={Accepted} skipped={Skipped} payload={Payload}", accepted, skipped, JsonSerializer.Serialize(new { accepted, skipped, scopes = scopes.Count })); } }