ModuleRebuildWorker.cs 6.0 KB

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