using Jiaowu.Api.Domain.Academic; using Jiaowu.Api.Domain.System; using Jiaowu.Api.Infrastructure.BackgroundJobs; using Jiaowu.Api.Infrastructure.Persistence; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; namespace Jiaowu.Api.Tests; public sealed class BackgroundJobOutboxTests { [Fact] public async Task In_memory_transport_keeps_job_types_on_independent_channels() { var transport = new InMemoryBackgroundJobTransport(); var automatic = new BackgroundJobEnvelope( Guid.NewGuid(), BackgroundJobKind.AutomaticSchedule, Guid.NewGuid()); var publish = new BackgroundJobEnvelope( Guid.NewGuid(), BackgroundJobKind.SchedulePublish, Guid.NewGuid()); await transport.PublishAsync(automatic, CancellationToken.None); await transport.PublishAsync(publish, CancellationToken.None); using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(1)); await using var automaticReader = transport .ReadAllAsync(BackgroundJobKind.AutomaticSchedule, timeout.Token) .GetAsyncEnumerator(timeout.Token); await using var publishReader = transport .ReadAllAsync(BackgroundJobKind.SchedulePublish, timeout.Token) .GetAsyncEnumerator(timeout.Token); Assert.True(await automaticReader.MoveNextAsync()); Assert.Equal(automatic, automaticReader.Current); Assert.True(await publishReader.MoveNextAsync()); Assert.Equal(publish, publishReader.Current); } [Fact] public async Task Publisher_recovers_pending_message_and_cleans_old_completion() { var databasePath = Path.Combine( Path.GetTempPath(), $"jiaowu-outbox-{Guid.NewGuid():N}.sqlite"); var services = new ServiceCollection(); services.AddLogging(); services.AddDbContext(options => options.UseSqlite($"Data Source={databasePath};Pooling=False")); var options = new BackgroundJobOptions { PollIntervalMilliseconds = 100, CompletedRetentionDays = 1, CleanupBatchSize = 10 }; services.AddSingleton(options); services.AddSingleton(); services.AddSingleton(provider => provider.GetRequiredService()); services.AddSingleton(); services.AddSingleton(); var provider = services.BuildServiceProvider(); try { Guid jobId; Guid oldCompletedOutboxId; await using (var scope = provider.CreateAsyncScope()) { var db = scope.ServiceProvider.GetRequiredService(); await db.Database.EnsureCreatedAsync(); var term = new AcademicTerm { Code = "RECOVERY", Name = "补投测试学期", AcademicYear = "2026-2027", Season = TermSeason.Autumn, StartDate = new DateOnly(2026, 9, 1), EndDate = new DateOnly(2027, 1, 15) }; var plan = new SchedulePlan { AcademicTerm = term, Name = "补投测试排课", Version = "V1" }; var job = new AutomaticScheduleJob { SchedulePlan = plan, ActiveSchedulePlanId = plan.Id }; var oldCompletedOutbox = BackgroundJobOutboxMessage.Create( BackgroundJobKind.SchedulePublish, Guid.NewGuid()); oldCompletedOutbox.State = BackgroundJobOutboxState.Completed; oldCompletedOutbox.CompletedAt = DateTime.UtcNow.AddDays(-2); oldCompletedOutboxId = oldCompletedOutbox.Id; jobId = job.Id; db.AddRange(term, plan, job, oldCompletedOutbox); await db.SaveChangesAsync(); } var publisher = provider.GetRequiredService(); await publisher.StartAsync(CancellationToken.None); try { var published = false; for (var attempt = 0; attempt < 50 && !published; attempt++) { await Task.Delay(50); await using var scope = provider.CreateAsyncScope(); var db = scope.ServiceProvider.GetRequiredService(); published = await db.BackgroundJobOutboxMessages.AsNoTracking() .AnyAsync(x => x.JobId == jobId && x.JobKind == BackgroundJobKind.AutomaticSchedule && x.State == BackgroundJobOutboxState.Published) && !await db.BackgroundJobOutboxMessages.AsNoTracking() .AnyAsync(x => x.Id == oldCompletedOutboxId); } Assert.True(published); } finally { await publisher.StopAsync(CancellationToken.None); } } finally { await provider.DisposeAsync(); File.Delete(databasePath); } } [Fact] public async Task Monitoring_snapshot_reports_backlog_and_expired_leases() { var databasePath = Path.Combine( Path.GetTempPath(), $"jiaowu-monitoring-{Guid.NewGuid():N}.sqlite"); var options = new DbContextOptionsBuilder() .UseSqlite($"Data Source={databasePath};Pooling=False") .Options; try { await using var db = new AppDbContext(options); await db.Database.EnsureCreatedAsync(); var pending = BackgroundJobOutboxMessage.Create( BackgroundJobKind.AutomaticSchedule, Guid.NewGuid()); pending.CreatedAt = DateTime.UtcNow.AddMinutes(-2); var processing = BackgroundJobOutboxMessage.Create( BackgroundJobKind.MakeupExamAuto, Guid.NewGuid()); processing.State = BackgroundJobOutboxState.Processing; processing.LeaseExpiresAt = DateTime.UtcNow.AddMinutes(-1); db.AddRange(pending, processing); await db.SaveChangesAsync(); var snapshot = await new BackgroundJobMonitoringService(db) .GetSnapshotAsync(CancellationToken.None); Assert.Equal(1, snapshot.Pending); Assert.Equal(1, snapshot.Processing); Assert.Equal(1, snapshot.ExpiredLeases); Assert.NotNull(snapshot.OldestUnfinishedAt); Assert.True(snapshot.OldestUnfinishedAgeSeconds >= 100); } finally { File.Delete(databasePath); } } [Fact] public async Task Runner_stops_after_retry_limit_and_ignores_duplicate_delivery() { var databasePath = Path.Combine( Path.GetTempPath(), $"jiaowu-runner-{Guid.NewGuid():N}.sqlite"); var services = new ServiceCollection(); services.AddLogging(); services.AddDbContext(options => options.UseSqlite($"Data Source={databasePath};Pooling=False")); services.AddScoped(); services.AddScoped(); services.AddSingleton(new BackgroundJobOptions()); services.AddSingleton(); services.AddSingleton(provider => provider.GetRequiredService()); services.AddSingleton(); services.AddSingleton(); var provider = services.BuildServiceProvider(); try { BackgroundJobEnvelope envelope; await using (var scope = provider.CreateAsyncScope()) { var db = scope.ServiceProvider.GetRequiredService(); await db.Database.EnsureCreatedAsync(); var term = new AcademicTerm { Code = "OUTBOX", Name = "消息测试学期", AcademicYear = "2026-2027", Season = TermSeason.Autumn, StartDate = new DateOnly(2026, 9, 1), EndDate = new DateOnly(2027, 1, 15) }; var plan = new SchedulePlan { AcademicTerm = term, Name = "消息测试排课", Version = "V1" }; var job = new AutomaticScheduleJob { SchedulePlan = plan, ActiveSchedulePlanId = plan.Id, Status = AutomaticScheduleJobStatus.Queued }; var outbox = BackgroundJobOutboxMessage.Create( BackgroundJobKind.AutomaticSchedule, job.Id); outbox.State = BackgroundJobOutboxState.Published; outbox.ProcessingAttempts = 5; db.AddRange(term, plan, job, outbox); await db.SaveChangesAsync(); envelope = new BackgroundJobEnvelope( outbox.Id, outbox.JobKind, outbox.JobId); } var runner = provider.GetRequiredService(); var first = await runner.RunAsync(envelope, CancellationToken.None); var duplicate = await runner.RunAsync(envelope, CancellationToken.None); Assert.Equal(BackgroundJobRunOutcome.Completed, first.Outcome); Assert.Equal(BackgroundJobRunOutcome.Completed, duplicate.Outcome); await using var assertScope = provider.CreateAsyncScope(); var assertDb = assertScope.ServiceProvider.GetRequiredService(); var persistedOutbox = await assertDb.BackgroundJobOutboxMessages.SingleAsync(); Assert.Equal(BackgroundJobOutboxState.Completed, persistedOutbox.State); Assert.NotNull(persistedOutbox.CompletedAt); Assert.Null(persistedOutbox.ProcessingToken); Assert.Null(persistedOutbox.LeaseExpiresAt); Assert.Equal(6, persistedOutbox.ProcessingAttempts); var persistedJob = await assertDb.AutomaticScheduleJobs.SingleAsync(); Assert.Equal(AutomaticScheduleJobStatus.Failed, persistedJob.Status); Assert.Null(persistedJob.ActiveSchedulePlanId); Assert.Contains("超过 5 次", persistedJob.ErrorMessage); } finally { await provider.DisposeAsync(); File.Delete(databasePath); } } }