using System.Diagnostics; using Jiaowu.Api.Domain.Academic; using Jiaowu.Api.Domain.System; using Jiaowu.Api.Infrastructure.Exams; using Jiaowu.Api.Infrastructure.Persistence; using Jiaowu.Api.Infrastructure.Scheduling; using Microsoft.EntityFrameworkCore; namespace Jiaowu.Api.Infrastructure.BackgroundJobs; public sealed record BackgroundJobRunResult( BackgroundJobRunOutcome Outcome, TimeSpan RetryAfter) { public static BackgroundJobRunResult Completed { get; } = new(BackgroundJobRunOutcome.Completed, TimeSpan.Zero); public static BackgroundJobRunResult Retry(TimeSpan retryAfter) => new(BackgroundJobRunOutcome.Retry, retryAfter); } public enum BackgroundJobRunOutcome { Completed, Retry } public sealed class BackgroundJobRunner( IServiceScopeFactory scopeFactory, IBackgroundJobTransport transport, BackgroundJobOptions options, BackgroundJobTelemetry telemetry, ILogger logger) { public async Task RunAsync( BackgroundJobEnvelope message, CancellationToken cancellationToken) { var startedAt = Stopwatch.GetTimestamp(); var token = Guid.NewGuid(); var claim = await TryClaimAsync(message, token, cancellationToken); if (!claim.Claimed) { return claim.Completed ? BackgroundJobRunResult.Completed : BackgroundJobRunResult.Retry(claim.RetryAfter); } if (claim.AttemptsExceeded) { await MarkJobRetryLimitExceededAsync(message, cancellationToken); await MarkCompletedAsync(message.OutboxMessageId, token, cancellationToken); telemetry.RecordProcessing( message.JobKind, "retry-limit", Stopwatch.GetElapsedTime(startedAt)); return BackgroundJobRunResult.Completed; } using var heartbeatCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); var heartbeat = RunHeartbeatAsync( message.OutboxMessageId, token, heartbeatCancellation.Token); try { await using var scope = scopeFactory.CreateAsyncScope(); switch (message.JobKind) { case BackgroundJobKind.AutomaticSchedule: await scope.ServiceProvider .GetRequiredService() .ProcessAsync(message.JobId, cancellationToken); break; case BackgroundJobKind.SchedulePublish: await scope.ServiceProvider .GetRequiredService() .ProcessAsync(message.JobId, cancellationToken); break; case BackgroundJobKind.MakeupExamAuto: await scope.ServiceProvider .GetRequiredService() .ProcessAsync(message.JobId, cancellationToken); break; case BackgroundJobKind.ExamArrangement: await scope.ServiceProvider .GetRequiredService() .ProcessAsync(message.JobId, cancellationToken); break; case BackgroundJobKind.ExamSignInExport: await scope.ServiceProvider .GetRequiredService() .ProcessAsync(message.JobId, cancellationToken); break; default: throw new InvalidOperationException( $"Unsupported background job kind '{message.JobKind}'."); } await MarkCompletedAsync(message.OutboxMessageId, token, cancellationToken); telemetry.RecordProcessing( message.JobKind, "completed", Stopwatch.GetElapsedTime(startedAt)); return BackgroundJobRunResult.Completed; } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { throw; } catch (Exception exception) { logger.LogError( exception, "Background job {JobKind}/{JobId} will be retried.", message.JobKind, message.JobId); await ReleaseForRetryAsync( message.OutboxMessageId, token, exception, CancellationToken.None); telemetry.RecordProcessing( message.JobKind, "retry", Stopwatch.GetElapsedTime(startedAt)); return BackgroundJobRunResult.Retry(TimeSpan.FromSeconds(2)); } finally { await heartbeatCancellation.CancelAsync(); try { await heartbeat; } catch (OperationCanceledException) { // Expected when processing completes or the application stops. } } } private async Task<( bool Claimed, bool Completed, bool AttemptsExceeded, TimeSpan RetryAfter)> TryClaimAsync( BackgroundJobEnvelope message, Guid token, CancellationToken cancellationToken) { await using var scope = scopeFactory.CreateAsyncScope(); var db = scope.ServiceProvider.GetRequiredService(); var now = DateTime.UtcNow; var leaseExpiresAt = now.AddSeconds(options.LeaseSeconds); var claimed = await db.BackgroundJobOutboxMessages .Where(x => x.Id == message.OutboxMessageId && x.JobKind == message.JobKind && x.JobId == message.JobId && (x.State == BackgroundJobOutboxState.Pending || x.State == BackgroundJobOutboxState.Publishing || x.State == BackgroundJobOutboxState.Published || (x.State == BackgroundJobOutboxState.Processing && x.LeaseExpiresAt != null && x.LeaseExpiresAt < now))) .ExecuteUpdateAsync( setters => setters .SetProperty(x => x.State, BackgroundJobOutboxState.Processing) .SetProperty(x => x.ProcessingToken, token) .SetProperty(x => x.LeaseExpiresAt, leaseExpiresAt) .SetProperty( x => x.ProcessingAttempts, x => x.ProcessingAttempts + 1) .SetProperty(x => x.LastError, (string?)null), cancellationToken); if (claimed == 1) { var attempts = await db.BackgroundJobOutboxMessages.AsNoTracking() .Where(x => x.Id == message.OutboxMessageId) .Select(x => x.ProcessingAttempts) .SingleAsync(cancellationToken); return ( true, false, attempts > options.ProcessingAttemptLimit, TimeSpan.Zero); } var current = await db.BackgroundJobOutboxMessages.AsNoTracking() .Where(x => x.Id == message.OutboxMessageId) .Select(x => new { x.State, x.LeaseExpiresAt }) .FirstOrDefaultAsync(cancellationToken); if (current is null || current.State == BackgroundJobOutboxState.Completed) return (false, true, false, TimeSpan.Zero); var retryAfter = current.LeaseExpiresAt.HasValue ? current.LeaseExpiresAt.Value - now : TimeSpan.FromSeconds(2); retryAfter = TimeSpan.FromSeconds(Math.Clamp( retryAfter.TotalSeconds, 1, options.LeaseSeconds)); return (false, false, false, retryAfter); } private async Task MarkJobRetryLimitExceededAsync( BackgroundJobEnvelope message, CancellationToken cancellationToken) { await using var scope = scopeFactory.CreateAsyncScope(); var db = scope.ServiceProvider.GetRequiredService(); var completedAt = DateTime.UtcNow; var error = $"后台任务连续处理失败超过 {options.ProcessingAttemptLimit} 次," + "已停止自动重试,请检查服务日志后重新创建任务。"; switch (message.JobKind) { case BackgroundJobKind.AutomaticSchedule: await db.AutomaticScheduleJobs .Where(x => x.Id == message.JobId && x.Status != AutomaticScheduleJobStatus.Succeeded && x.Status != AutomaticScheduleJobStatus.Failed) .ExecuteUpdateAsync( setters => setters .SetProperty( x => x.Status, AutomaticScheduleJobStatus.Failed) .SetProperty(x => x.ActiveSchedulePlanId, (Guid?)null) .SetProperty(x => x.ErrorMessage, error) .SetProperty(x => x.CompletedAt, completedAt), cancellationToken); break; case BackgroundJobKind.SchedulePublish: await db.SchedulePublishJobs .Where(x => x.Id == message.JobId && x.Status != SchedulePublishJobStatus.Succeeded && x.Status != SchedulePublishJobStatus.Failed) .ExecuteUpdateAsync( setters => setters .SetProperty( x => x.Status, SchedulePublishJobStatus.Failed) .SetProperty(x => x.ActiveAcademicTermId, (Guid?)null) .SetProperty(x => x.CurrentStep, "后台处理已停止") .SetProperty(x => x.ErrorMessage, error) .SetProperty(x => x.CompletedAt, completedAt), cancellationToken); break; case BackgroundJobKind.MakeupExamAuto: await db.MakeupExamAutoJobs .Where(x => x.Id == message.JobId && x.Status != MakeupExamAutoJobStatus.Succeeded && x.Status != MakeupExamAutoJobStatus.Failed) .ExecuteUpdateAsync( setters => setters .SetProperty( x => x.Status, MakeupExamAutoJobStatus.Failed) .SetProperty(x => x.ErrorMessage, error) .SetProperty(x => x.CompletedAt, completedAt), cancellationToken); break; case BackgroundJobKind.ExamArrangement: await db.ExamArrangementJobs .Where(x => x.Id == message.JobId && x.Status != ExamArrangementJobStatus.Succeeded && x.Status != ExamArrangementJobStatus.Failed) .ExecuteUpdateAsync( setters => setters .SetProperty( x => x.Status, ExamArrangementJobStatus.Failed) .SetProperty(x => x.ActivePlanId, (Guid?)null) .SetProperty(x => x.CurrentStep, "后台处理已停止") .SetProperty(x => x.ErrorMessage, error) .SetProperty(x => x.CompletedAt, completedAt), cancellationToken); break; default: throw new ArgumentOutOfRangeException( nameof(message.JobKind), message.JobKind, null); } logger.LogError( "Background job {JobKind}/{JobId} exceeded {AttemptLimit} processing attempts.", message.JobKind, message.JobId, options.ProcessingAttemptLimit); } private async Task RunHeartbeatAsync( Guid outboxMessageId, Guid token, CancellationToken cancellationToken) { var intervalSeconds = Math.Max(5, options.LeaseSeconds / 3); while (!cancellationToken.IsCancellationRequested) { await Task.Delay(TimeSpan.FromSeconds(intervalSeconds), cancellationToken); await using var scope = scopeFactory.CreateAsyncScope(); var db = scope.ServiceProvider.GetRequiredService(); await db.BackgroundJobOutboxMessages .Where(x => x.Id == outboxMessageId && x.State == BackgroundJobOutboxState.Processing && x.ProcessingToken == token) .ExecuteUpdateAsync( setters => setters.SetProperty( x => x.LeaseExpiresAt, DateTime.UtcNow.AddSeconds(options.LeaseSeconds)), cancellationToken); } } private async Task MarkCompletedAsync( Guid outboxMessageId, Guid token, CancellationToken cancellationToken) { await using var scope = scopeFactory.CreateAsyncScope(); var db = scope.ServiceProvider.GetRequiredService(); await db.BackgroundJobOutboxMessages .Where(x => x.Id == outboxMessageId && x.State == BackgroundJobOutboxState.Processing && x.ProcessingToken == token) .ExecuteUpdateAsync( setters => setters .SetProperty(x => x.State, BackgroundJobOutboxState.Completed) .SetProperty(x => x.CompletedAt, DateTime.UtcNow) .SetProperty(x => x.ProcessingToken, (Guid?)null) .SetProperty(x => x.LeaseExpiresAt, (DateTime?)null), cancellationToken); } private async Task ReleaseForRetryAsync( Guid outboxMessageId, Guid token, Exception exception, CancellationToken cancellationToken) { await using var scope = scopeFactory.CreateAsyncScope(); var db = scope.ServiceProvider.GetRequiredService(); var message = exception.GetBaseException().Message; if (message.Length > 2000) message = message[..2000]; var state = transport.IsDurable ? BackgroundJobOutboxState.Published : BackgroundJobOutboxState.Pending; await db.BackgroundJobOutboxMessages .Where(x => x.Id == outboxMessageId && x.State == BackgroundJobOutboxState.Processing && x.ProcessingToken == token) .ExecuteUpdateAsync( setters => setters .SetProperty(x => x.State, state) .SetProperty(x => x.ProcessingToken, (Guid?)null) .SetProperty(x => x.LeaseExpiresAt, (DateTime?)null) .SetProperty(x => x.LastError, message), cancellationToken); } }