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 }));
}
}