From 865ce2f6d833f998bbb33f74a59bbd86c504050e Mon Sep 17 00:00:00 2001 From: iammukeshm Date: Tue, 29 Sep 2026 16:57:00 +0530 Subject: [PATCH] fix(auditing): purge audit records per tenant so the retention job runs AuditRetentionJob is a Hangfire recurring job registered without a tenant, so it ran with no Finbuckle context and the default-on tenant filter on AuditRecords threw a NullReferenceException on every run. Audit tables grew without bound in every tenant. The job now lists tenants from IMultiTenantStore, opens a scope per tenant, sets IMultiTenantContextSetter, and sweeps inside that context. A failure in one tenant is logged and does not stop the others. Adds an integration test that seeds old and recent records in two tenants and runs the job from a tenant-less scope. Closes #1411 Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Persistence/AuditRetentionJob.cs | 70 ++++++-- .../Tests/Auditing/AuditRetentionJobTests.cs | 167 ++++++++++++++++++ 2 files changed, 223 insertions(+), 14 deletions(-) create mode 100644 src/Tests/Integration.Tests/Tests/Auditing/AuditRetentionJobTests.cs diff --git a/src/Modules/Auditing/Modules.Auditing/Persistence/AuditRetentionJob.cs b/src/Modules/Auditing/Modules.Auditing/Persistence/AuditRetentionJob.cs index bba1c9bcc5..df86e372c4 100644 --- a/src/Modules/Auditing/Modules.Auditing/Persistence/AuditRetentionJob.cs +++ b/src/Modules/Auditing/Modules.Auditing/Persistence/AuditRetentionJob.cs @@ -1,30 +1,34 @@ +using Finbuckle.MultiTenant; +using Finbuckle.MultiTenant.Abstractions; +using FSH.Framework.Shared.Multitenancy; using FSH.Modules.Auditing.Contracts; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; namespace FSH.Modules.Auditing.Persistence; /// /// Daily Hangfire job that prunes the audit table per -/// . Uses ExecuteDeleteAsync with +/// , in every tenant. Uses ExecuteDeleteAsync with /// a bounded batch size so a single run doesn't take a long-held lock on /// the table — each event-type sweep loops until fewer than batch-size /// rows are deleted. /// public sealed class AuditRetentionJob { - private readonly AuditDbContext _db; + private readonly IServiceScopeFactory _scopeFactory; private readonly AuditRetentionOptions _opts; private readonly TimeProvider _timeProvider; private readonly ILogger _logger; public AuditRetentionJob( - AuditDbContext db, + IServiceScopeFactory scopeFactory, AuditRetentionOptions opts, TimeProvider timeProvider, ILogger logger) { - _db = db; + _scopeFactory = scopeFactory; _opts = opts; _timeProvider = timeProvider; _logger = logger; @@ -38,20 +42,58 @@ public async Task RunAsync(CancellationToken ct) return; } + // The recurring job is registered without a tenant, so none is resolved and the default-on + // tenant filter on AuditRecords has nothing to compare against (it throws). Each tenant is + // purged inside its own context, which also points AuditDbContext at a tenant's dedicated + // database when it has one. + List tenants; + using (var scope = _scopeFactory.CreateScope()) + { + var tenantStore = scope.ServiceProvider.GetRequiredService>(); + tenants = (await tenantStore.GetAllAsync().ConfigureAwait(false)).ToList(); + } + var now = _timeProvider.GetUtcNow().UtcDateTime; long total = 0; - total += await SweepAsync(AuditEventType.Activity, now.AddDays(-_opts.ActivityRetentionDays), ct).ConfigureAwait(false); - total += await SweepAsync(AuditEventType.EntityChange, now.AddDays(-_opts.EntityChangeRetentionDays), ct).ConfigureAwait(false); - total += await SweepAsync(AuditEventType.Security, now.AddDays(-_opts.SecurityRetentionDays), ct).ConfigureAwait(false); - total += await SweepAsync(AuditEventType.Exception, now.AddDays(-_opts.ExceptionRetentionDays), ct).ConfigureAwait(false); + foreach (var tenant in tenants) + { + ct.ThrowIfCancellationRequested(); + total += await PurgeTenantAsync(tenant, now, ct).ConfigureAwait(false); + } if (_logger.IsEnabled(LogLevel.Information)) { - _logger.LogInformation("[Auditing] retention job purged {Total} rows.", total); + _logger.LogInformation("[Auditing] retention job purged {Total} rows across {TenantCount} tenants.", + total, tenants.Count); + } + } + + private async Task PurgeTenantAsync(AppTenantInfo tenant, DateTime now, CancellationToken ct) + { + try + { + using var scope = _scopeFactory.CreateScope(); + scope.ServiceProvider.GetRequiredService() + .MultiTenantContext = new MultiTenantContext(tenant); + + var db = scope.ServiceProvider.GetRequiredService(); + long total = 0; + total += await SweepAsync(db, tenant.Id, AuditEventType.Activity, now.AddDays(-_opts.ActivityRetentionDays), ct).ConfigureAwait(false); + total += await SweepAsync(db, tenant.Id, AuditEventType.EntityChange, now.AddDays(-_opts.EntityChangeRetentionDays), ct).ConfigureAwait(false); + total += await SweepAsync(db, tenant.Id, AuditEventType.Security, now.AddDays(-_opts.SecurityRetentionDays), ct).ConfigureAwait(false); + total += await SweepAsync(db, tenant.Id, AuditEventType.Exception, now.AddDays(-_opts.ExceptionRetentionDays), ct).ConfigureAwait(false); + return total; + } + catch (Exception ex) when (ex is not OperationCanceledException) + { + // One tenant's database being unreachable must not keep the other tenants' audit rows around. + _logger.LogError(ex, "[Auditing] retention purge failed for tenant {TenantId}.", tenant.Id); + return 0; } } - private async Task SweepAsync(AuditEventType eventType, DateTime cutoffUtc, CancellationToken ct) + private async Task SweepAsync( + AuditDbContext db, string? tenantId, AuditEventType eventType, DateTime cutoffUtc, CancellationToken ct) { long swept = 0; var typeId = (int)eventType; @@ -61,10 +103,10 @@ private async Task SweepAsync(AuditEventType eventType, DateTime cutoffUtc { // Sub-query trick: ExecuteDeleteAsync doesn't support TOP/LIMIT // directly, so we filter to a bounded id-set first. - var deleted = await _db.AuditRecords + var deleted = await db.AuditRecords .Where(a => a.EventType == typeId && a.OccurredAtUtc < cutoffUtc - && _db.AuditRecords + && db.AuditRecords .Where(b => b.EventType == typeId && b.OccurredAtUtc < cutoffUtc) .OrderBy(b => b.OccurredAtUtc) .Select(b => b.Id) @@ -79,8 +121,8 @@ private async Task SweepAsync(AuditEventType eventType, DateTime cutoffUtc if (swept > 0 && _logger.IsEnabled(LogLevel.Information)) { - _logger.LogInformation("[Auditing] purged {Count} {EventType} events older than {Cutoff:o}.", - swept, eventType, cutoffUtc); + _logger.LogInformation("[Auditing] purged {Count} {EventType} events older than {Cutoff:o} for tenant {TenantId}.", + swept, eventType, cutoffUtc, tenantId); } return swept; } diff --git a/src/Tests/Integration.Tests/Tests/Auditing/AuditRetentionJobTests.cs b/src/Tests/Integration.Tests/Tests/Auditing/AuditRetentionJobTests.cs new file mode 100644 index 0000000000..2344bfffca --- /dev/null +++ b/src/Tests/Integration.Tests/Tests/Auditing/AuditRetentionJobTests.cs @@ -0,0 +1,167 @@ +using Finbuckle.MultiTenant; +using Finbuckle.MultiTenant.Abstractions; +using FSH.Framework.Shared.Multitenancy; +using FSH.Modules.Auditing; +using FSH.Modules.Auditing.Contracts; +using FSH.Modules.Auditing.Persistence; +using FSH.Modules.Multitenancy.Contracts.Dtos; +using Integration.Tests.Infrastructure; + +namespace Integration.Tests.Tests.Auditing; + +/// +/// The audit retention purge is a Hangfire recurring job registered without a tenant parameter, so it +/// runs with no ambient tenant. AuditRecords carry the default-on tenant filter, which means the purge +/// only reaches a tenant's rows when it runs inside that tenant's context. Records are seeded in two +/// tenants and the job is resolved from a fresh, tenant-less scope — the same shape Hangfire's +/// activator gives it — and run once. +/// +[Collection(FshCollectionDefinition.Name)] +public sealed class AuditRetentionJobTests +{ + // Activity events are kept for 30 days by default; 31 is safely past the cutoff, 1 is inside it. + private const int ActivityRetentionDays = 30; + private const int DaysPastRetention = 31; + private const int DaysInsideRetention = 1; + + private readonly FshWebApplicationFactory _factory; + private readonly AuthHelper _auth; + + public AuditRetentionJobTests(FshWebApplicationFactory factory) + { + _factory = factory; + _auth = new AuthHelper(factory); + } + + [Fact] + public async Task RunAsync_Should_PurgeRecordsPastRetention_InEveryTenant_And_KeepTheRest() + { + // Arrange + using var rootClient = await _auth.CreateRootAdminClientAsync(); + var uniqueId = Guid.NewGuid().ToString("N")[..8]; + var otherTenantId = $"audit-ret-{uniqueId}"; + var otherAdminEmail = $"audit-ret-admin-{uniqueId}@tenant.com"; + + await CreateTenantAsync(rootClient, otherTenantId, otherAdminEmail); + await WaitForProvisioningAsync(rootClient, otherTenantId); + + var now = DateTime.UtcNow; + var rootOld = await SeedAuditAsync(TestConstants.RootTenantId, now.AddDays(-DaysPastRetention)); + var otherOld = await SeedAuditAsync(otherTenantId, now.AddDays(-DaysPastRetention)); + var rootRecent = await SeedAuditAsync(TestConstants.RootTenantId, now.AddDays(-DaysInsideRetention)); + var otherRecent = await SeedAuditAsync(otherTenantId, now.AddDays(-DaysInsideRetention)); + + var options = new AuditRetentionOptions + { + Enabled = true, + ActivityRetentionDays = ActivityRetentionDays, + }; + + // Act — a fresh scope with no tenant set, as FshJobActivator gives a job registered via AddOrUpdate. + using (var jobScope = _factory.Services.CreateScope()) + { + var job = ActivatorUtilities.CreateInstance(jobScope.ServiceProvider, options); + await job.RunAsync(CancellationToken.None); + } + + // Assert + (await AuditExistsAsync(TestConstants.RootTenantId, rootOld)) + .ShouldBeFalse("a root-tenant audit record past retention must be purged"); + (await AuditExistsAsync(otherTenantId, otherOld)) + .ShouldBeFalse("an audit record past retention in a non-root tenant must be purged"); + (await AuditExistsAsync(TestConstants.RootTenantId, rootRecent)) + .ShouldBeTrue("a root-tenant audit record inside retention must be kept"); + (await AuditExistsAsync(otherTenantId, otherRecent)) + .ShouldBeTrue("an audit record inside retention in a non-root tenant must be kept"); + } + + // Tenant context is an AsyncLocal, so it is set in the same method as the DbContext call. + private async Task SeedAuditAsync(string tenantId, DateTime occurredAtUtc) + { + using var scope = _factory.Services.CreateScope(); + var tenant = await scope.ServiceProvider + .GetRequiredService>() + .GetAsync(tenantId); + tenant.ShouldNotBeNull(); + scope.ServiceProvider.GetRequiredService() + .MultiTenantContext = new MultiTenantContext(tenant); + + var record = new AuditRecord + { + Id = Guid.NewGuid(), + OccurredAtUtc = occurredAtUtc, + ReceivedAtUtc = occurredAtUtc, + EventType = (int)AuditEventType.Activity, + Severity = (byte)AuditSeverity.Information, + TenantId = tenantId, + Source = "audit-retention-test", + PayloadJson = "{}", + }; + + var db = scope.ServiceProvider.GetRequiredService(); + db.AuditRecords.Add(record); + await db.SaveChangesAsync(); + return record.Id; + } + + private async Task AuditExistsAsync(string tenantId, Guid auditId) + { + using var scope = _factory.Services.CreateScope(); + var tenant = await scope.ServiceProvider + .GetRequiredService>() + .GetAsync(tenantId); + tenant.ShouldNotBeNull(); + scope.ServiceProvider.GetRequiredService() + .MultiTenantContext = new MultiTenantContext(tenant); + + var db = scope.ServiceProvider.GetRequiredService(); + return await db.AuditRecords.AsNoTracking().AnyAsync(a => a.Id == auditId); + } + + private static async Task CreateTenantAsync(HttpClient rootClient, string tenantId, string adminEmail) + { + var response = await rootClient.PostAsJsonAsync(TestConstants.TenantsBasePath, new + { + id = tenantId, + name = $"Tenant {tenantId}", + connectionString = (string?)null, + adminEmail, + adminPassword = TestConstants.DefaultPassword, + issuer = $"{tenantId}.issuer" + }); + var body = await response.Content.ReadAsStringAsync(); + response.StatusCode.ShouldBe(HttpStatusCode.Created, $"Create tenant failed: {body}"); + } + + // The status body also lists each step, and a finished step reads "Completed" while later steps + // are still running, so only the overall Status field is trusted. + private static async Task WaitForProvisioningAsync(HttpClient client, string tenantId) + { + const int maxRetries = 60; + for (int i = 0; i < maxRetries; i++) + { + var statusResponse = await client.GetAsync( + $"{TestConstants.TenantsBasePath}/{tenantId}/provisioning"); + + if (statusResponse.IsSuccessStatusCode) + { + var status = await statusResponse.Content.ReadFromJsonAsync(); + if (string.Equals(status?.Status, "Completed", StringComparison.Ordinal)) + { + return; + } + + if (string.Equals(status?.Status, "Failed", StringComparison.Ordinal)) + { + throw new InvalidOperationException( + $"Tenant {tenantId} provisioning failed at {status?.CurrentStep}: {status?.Error}"); + } + } + + await Task.Delay(1000); + } + + throw new TimeoutException( + $"Tenant {tenantId} provisioning did not complete within {maxRetries} seconds."); + } +}