ModuleRebuildWorker.cs 5.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151
  1. using Microsoft.Extensions.DependencyInjection;
  2. using Microsoft.Extensions.Hosting;
  3. using Microsoft.Extensions.Logging;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  5. public sealed class ModuleRebuildWorker : BackgroundService
  6. {
  7. private static readonly TimeSpan IdleDelay = TimeSpan.FromSeconds(5);
  8. private readonly IServiceScopeFactory _scopeFactory;
  9. private readonly ModuleRebuildQueue _queue;
  10. private readonly ILogger _logger;
  11. public ModuleRebuildWorker(
  12. IServiceScopeFactory scopeFactory,
  13. ModuleRebuildQueue queue,
  14. ILoggerFactory loggerFactory)
  15. {
  16. _scopeFactory = scopeFactory;
  17. _queue = queue;
  18. _logger = loggerFactory.CreateLogger(nameof(ModuleRebuildWorker));
  19. }
  20. protected override async Task ExecuteAsync(CancellationToken stoppingToken)
  21. {
  22. try { await Task.Delay(TimeSpan.FromSeconds(8), stoppingToken); }
  23. catch (OperationCanceledException) { return; }
  24. await FailStaleAsync(stoppingToken);
  25. var running = new List<Task>();
  26. while (!stoppingToken.IsCancellationRequested)
  27. {
  28. try
  29. {
  30. await FillSlotsAsync(running, stoppingToken);
  31. }
  32. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  33. {
  34. break;
  35. }
  36. catch (Exception ex)
  37. {
  38. _logger.LogWarning(ex, "[ModuleRebuildWorker] run failed");
  39. }
  40. running.RemoveAll(t => t.IsCompleted);
  41. if (running.Count == 0)
  42. {
  43. try
  44. {
  45. using var linked = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
  46. linked.CancelAfter(IdleDelay);
  47. try { await _queue.Reader.ReadAsync(linked.Token); }
  48. catch (OperationCanceledException) when (!stoppingToken.IsCancellationRequested) { }
  49. }
  50. catch (OperationCanceledException) { break; }
  51. }
  52. else
  53. {
  54. var pulse = WaitPulseAsync(stoppingToken);
  55. var completed = await Task.WhenAny(running.Append(pulse));
  56. if (completed != pulse)
  57. running.Remove(completed);
  58. }
  59. }
  60. try { await Task.WhenAll(running); }
  61. catch (Exception ex) { _logger.LogWarning(ex, "[ModuleRebuildWorker] drain failed"); }
  62. }
  63. private async Task FillSlotsAsync(List<Task> running, CancellationToken stoppingToken)
  64. {
  65. int max = 2;
  66. try
  67. {
  68. using var probe = _scopeFactory.CreateScope();
  69. max = probe.ServiceProvider.GetRequiredService<IModuleRebuildCapability>().MaxParallelScopes;
  70. }
  71. catch
  72. {
  73. max = 2;
  74. }
  75. while (running.Count < max && !stoppingToken.IsCancellationRequested)
  76. {
  77. var started = await TryStartOneAsync(stoppingToken);
  78. if (started == null)
  79. break;
  80. running.Add(started);
  81. }
  82. }
  83. private async Task<Task> TryStartOneAsync(CancellationToken stoppingToken)
  84. {
  85. using var probe = _scopeFactory.CreateScope();
  86. var capability = probe.ServiceProvider.GetRequiredService<IModuleRebuildCapability>();
  87. var enabled = MdpRebuildScope.RebuildModules.Where(capability.IsEnabled).ToArray();
  88. if (enabled.Length == 0)
  89. return null;
  90. var job = await probe.ServiceProvider.GetRequiredService<IModuleRebuildJobStore>()
  91. .ClaimNextQueuedAsync(enabled, stoppingToken);
  92. if (job == null)
  93. return null;
  94. return ExecuteClaimedAsync(job, stoppingToken);
  95. }
  96. private async Task ExecuteClaimedAsync(Admin.NET.Plugin.AiDOP.Entity.SmartOps.AdoModuleDashboardRebuildJob job, CancellationToken stoppingToken)
  97. {
  98. using var scope = _scopeFactory.CreateScope();
  99. var svc = scope.ServiceProvider.GetRequiredService<ModuleRebuildService>();
  100. var runLock = scope.ServiceProvider.GetRequiredService<IModuleRebuildLock>();
  101. var handler = scope.ServiceProvider.GetServices<IModuleRebuildHandler>()
  102. .FirstOrDefault(x => string.Equals(x.ModuleCode, job.ModuleCode, StringComparison.OrdinalIgnoreCase));
  103. if (handler == null)
  104. {
  105. job.Status = ModuleRebuildStatus.Failed;
  106. job.CurrentStage = ModuleRebuildStages.Failed;
  107. job.ErrorMessage = $"未注册 {job.ModuleCode} 重算处理器";
  108. job.FinishedAt = DateTime.Now;
  109. job.UpdateTime = DateTime.Now;
  110. await scope.ServiceProvider.GetRequiredService<IModuleRebuildJobStore>().UpdateAsync(job, CancellationToken.None);
  111. return;
  112. }
  113. // 进度写入已带 heartbeat;后台另开连接会与单例 ISqlSugarClient 抢同一条 MySQL 连接。
  114. await svc.RunClaimedAsync(job, handler, runLock, stoppingToken);
  115. }
  116. private async Task WaitPulseAsync(CancellationToken stoppingToken)
  117. {
  118. using var linked = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
  119. linked.CancelAfter(IdleDelay);
  120. try { await _queue.Reader.ReadAsync(linked.Token); }
  121. catch (OperationCanceledException) { }
  122. }
  123. private async Task FailStaleAsync(CancellationToken stoppingToken)
  124. {
  125. try
  126. {
  127. using var scope = _scopeFactory.CreateScope();
  128. await scope.ServiceProvider.GetRequiredService<ModuleRebuildService>().FailStaleAsync(stoppingToken);
  129. }
  130. catch (Exception ex)
  131. {
  132. _logger.LogWarning(ex, "[ModuleRebuildWorker] fail-stale failed");
  133. }
  134. }
  135. }