Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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;

/// <summary>
/// Daily Hangfire job that prunes the audit table per
/// <see cref="AuditRetentionOptions"/>. Uses <c>ExecuteDeleteAsync</c> with
/// <see cref="AuditRetentionOptions"/>, in every tenant. Uses <c>ExecuteDeleteAsync</c> 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.
/// </summary>
public sealed class AuditRetentionJob
{
private readonly AuditDbContext _db;
private readonly IServiceScopeFactory _scopeFactory;
private readonly AuditRetentionOptions _opts;
private readonly TimeProvider _timeProvider;
private readonly ILogger<AuditRetentionJob> _logger;

public AuditRetentionJob(
AuditDbContext db,
IServiceScopeFactory scopeFactory,
AuditRetentionOptions opts,
TimeProvider timeProvider,
ILogger<AuditRetentionJob> logger)
{
_db = db;
_scopeFactory = scopeFactory;
_opts = opts;
_timeProvider = timeProvider;
_logger = logger;
Expand All @@ -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<AppTenantInfo> tenants;
using (var scope = _scopeFactory.CreateScope())
{
var tenantStore = scope.ServiceProvider.GetRequiredService<IMultiTenantStore<AppTenantInfo>>();
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<long> PurgeTenantAsync(AppTenantInfo tenant, DateTime now, CancellationToken ct)
{
try
{
using var scope = _scopeFactory.CreateScope();
scope.ServiceProvider.GetRequiredService<IMultiTenantContextSetter>()
.MultiTenantContext = new MultiTenantContext<AppTenantInfo>(tenant);

var db = scope.ServiceProvider.GetRequiredService<AuditDbContext>();
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<long> SweepAsync(AuditEventType eventType, DateTime cutoffUtc, CancellationToken ct)
private async Task<long> SweepAsync(
AuditDbContext db, string? tenantId, AuditEventType eventType, DateTime cutoffUtc, CancellationToken ct)
{
long swept = 0;
var typeId = (int)eventType;
Expand All @@ -61,10 +103,10 @@ private async Task<long> 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)
Expand All @@ -79,8 +121,8 @@ private async Task<long> 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;
}
Expand Down
167 changes: 167 additions & 0 deletions src/Tests/Integration.Tests/Tests/Auditing/AuditRetentionJobTests.cs
Original file line number Diff line number Diff line change
@@ -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;

/// <summary>
/// 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.
/// </summary>
[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<AuditRetentionJob>(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<Guid> SeedAuditAsync(string tenantId, DateTime occurredAtUtc)
{
using var scope = _factory.Services.CreateScope();
var tenant = await scope.ServiceProvider
.GetRequiredService<IMultiTenantStore<AppTenantInfo>>()
.GetAsync(tenantId);
tenant.ShouldNotBeNull();
scope.ServiceProvider.GetRequiredService<IMultiTenantContextSetter>()
.MultiTenantContext = new MultiTenantContext<AppTenantInfo>(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<AuditDbContext>();
db.AuditRecords.Add(record);
await db.SaveChangesAsync();
return record.Id;
}

private async Task<bool> AuditExistsAsync(string tenantId, Guid auditId)
{
using var scope = _factory.Services.CreateScope();
var tenant = await scope.ServiceProvider
.GetRequiredService<IMultiTenantStore<AppTenantInfo>>()
.GetAsync(tenantId);
tenant.ShouldNotBeNull();
scope.ServiceProvider.GetRequiredService<IMultiTenantContextSetter>()
.MultiTenantContext = new MultiTenantContext<AppTenantInfo>(tenant);

var db = scope.ServiceProvider.GetRequiredService<AuditDbContext>();
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<TenantProvisioningStatusDto>();
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.");
}
}
Loading