diff --git a/src/Exceptionless.Core/Jobs/CleanupOrphanedDataJob.cs b/src/Exceptionless.Core/Jobs/CleanupOrphanedDataJob.cs index 37a67ae24a..d6967266d7 100644 --- a/src/Exceptionless.Core/Jobs/CleanupOrphanedDataJob.cs +++ b/src/Exceptionless.Core/Jobs/CleanupOrphanedDataJob.cs @@ -23,6 +23,7 @@ namespace Exceptionless.Core.Jobs; [Job(Description = "Deletes orphaned data.", IsContinuous = false)] public class CleanupOrphanedDataJob : JobWithLockBase, IHealthCheck { + internal static readonly TimeSpan OrphanedEventLookback = TimeSpan.FromDays(3); private readonly ExceptionlessElasticConfiguration _config; private readonly ElasticsearchClient _elasticClient; private readonly IStackRepository _stackRepository; @@ -58,21 +59,30 @@ ILoggerFactory loggerFactory protected override async Task RunInternalAsync(JobContext context) { - await DeleteOrphanedEventsByStackAsync(context); - await DeleteOrphanedEventsByProjectAsync(context); - await DeleteOrphanedEventsByOrganizationAsync(context); + _lastRun = _timeProvider.GetUtcNow().UtcDateTime; + + var orphanedEventCutoffUtc = GetOrphanedEventCutoffUtc(); + await DeleteOrphanedEventsByStackAsync(context, orphanedEventCutoffUtc); + await DeleteOrphanedEventsByProjectAsync(context, orphanedEventCutoffUtc); + await DeleteOrphanedEventsByOrganizationAsync(context, orphanedEventCutoffUtc); await FixDuplicateStacks(context); return JobResult.Success; } - public async Task DeleteOrphanedEventsByStackAsync(JobContext context) + public Task DeleteOrphanedEventsByStackAsync(JobContext context) + { + return DeleteOrphanedEventsByStackAsync(context, GetOrphanedEventCutoffUtc()); + } + + private async Task DeleteOrphanedEventsByStackAsync(JobContext context, DateTime orphanedEventCutoffUtc) { // get approximate number of unique stack ids var stackCardinality = await _elasticClient.SearchAsync(s => s .Indices(GetEventIndexPattern()) .Size(0) + .Query(q => RecentEventQuery(q, orphanedEventCutoffUtc)) .AddAggregation("cardinality_stack_id", a => a.Cardinality(c => c.Field(f => f.StackId).PrecisionThreshold(40000)))); double? uniqueStackIdCount = stackCardinality.Aggregations?.GetCardinality("cardinality_stack_id")?.Value; @@ -93,6 +103,7 @@ public async Task DeleteOrphanedEventsByStackAsync(JobContext context) var stackIdTerms = await _elasticClient.SearchAsync(s => s .Indices(GetEventIndexPattern()) .Size(0) + .Query(q => RecentEventQuery(q, orphanedEventCutoffUtc)) .AddAggregation("terms_stack_id", a => a.Terms(c => c.Field(f => f.StackId).Include(new TermsInclude(batchNumber, buckets)).Size(batchSize * 2)))); string[] stackIds = stackIdTerms.Aggregations?.GetStringTerms("terms_stack_id")?.Buckets.Select(b => b.Key.ToString()!).ToArray() ?? []; @@ -115,18 +126,26 @@ public async Task DeleteOrphanedEventsByStackAsync(JobContext context) _logger.LogInformation("{BatchNumber}/{BatchCount}: Found {OrphanedEventCount} orphaned events from missing stacks {MissingStackIds} out of {StackIdCount}", batchNumber, buckets, missingStackIds.Length, missingStackIds, stackIds.Length); await _elasticClient.DeleteByQueryAsync(r => r .Indices(GetEventIndexPattern()) - .Query(q => q.Terms(t => t.Field(f => f.StackId).Terms(new TermsQueryField(missingStackIds.Select(FieldValueHelper.ToFieldValue).ToList()))))); + .Query(q => q.Bool(b => b.Filter( + f => f.Terms(t => t.Field(e => e.StackId).Terms(new TermsQueryField(missingStackIds.Select(FieldValueHelper.ToFieldValue).ToList()))), + f => RecentEventQuery(f, orphanedEventCutoffUtc))))); } - _logger.LogInformation("Found {OrphanedEventCount} orphaned events from missing stacks out of {StackIdCount}", totalOrphanedEventCount, totalStackIds); + _logger.LogInformation("Found {OrphanedEventCount} orphaned events from missing stacks out of {StackIdCount} since {OrphanedEventCutoffUtc}", totalOrphanedEventCount, totalStackIds, orphanedEventCutoffUtc); } - public async Task DeleteOrphanedEventsByProjectAsync(JobContext context) + public Task DeleteOrphanedEventsByProjectAsync(JobContext context) + { + return DeleteOrphanedEventsByProjectAsync(context, GetOrphanedEventCutoffUtc()); + } + + private async Task DeleteOrphanedEventsByProjectAsync(JobContext context, DateTime orphanedEventCutoffUtc) { // get approximate number of unique project ids var projectCardinality = await _elasticClient.SearchAsync(s => s .Indices(GetEventIndexPattern()) .Size(0) + .Query(q => RecentEventQuery(q, orphanedEventCutoffUtc)) .AddAggregation("cardinality_project_id", a => a.Cardinality(c => c.Field(f => f.ProjectId).PrecisionThreshold(40000)))); double? uniqueProjectIdCount = projectCardinality.Aggregations?.GetCardinality("cardinality_project_id")?.Value; @@ -147,6 +166,7 @@ public async Task DeleteOrphanedEventsByProjectAsync(JobContext context) var projectIdTerms = await _elasticClient.SearchAsync(s => s .Indices(GetEventIndexPattern()) .Size(0) + .Query(q => RecentEventQuery(q, orphanedEventCutoffUtc)) .AddAggregation("terms_project_id", a => a.Terms(c => c.Field(f => f.ProjectId).Include(new TermsInclude(batchNumber, buckets)).Size(batchSize * 2)))); string[] projectIds = projectIdTerms.Aggregations?.GetStringTerms("terms_project_id")?.Buckets.Select(b => b.Key.ToString()!).ToArray() ?? []; @@ -169,18 +189,26 @@ public async Task DeleteOrphanedEventsByProjectAsync(JobContext context) _logger.LogInformation("{BatchNumber}/{BatchCount}: Found {OrphanedEventCount} orphaned events from missing projects {MissingProjectIds} out of {ProjectIdCount}", batchNumber, buckets, missingProjectIds.Length, missingProjectIds, projectIds.Length); await _elasticClient.DeleteByQueryAsync(r => r .Indices(GetEventIndexPattern()) - .Query(q => q.Terms(t => t.Field(f => f.ProjectId).Terms(new TermsQueryField(missingProjectIds.Select(FieldValueHelper.ToFieldValue).ToList()))))); + .Query(q => q.Bool(b => b.Filter( + f => f.Terms(t => t.Field(e => e.ProjectId).Terms(new TermsQueryField(missingProjectIds.Select(FieldValueHelper.ToFieldValue).ToList()))), + f => RecentEventQuery(f, orphanedEventCutoffUtc))))); } - _logger.LogInformation("Found {OrphanedEventCount} orphaned events from missing projects out of {ProjectIdCount}", totalOrphanedEventCount, totalProjectIds); + _logger.LogInformation("Found {OrphanedEventCount} orphaned events from missing projects out of {ProjectIdCount} since {OrphanedEventCutoffUtc}", totalOrphanedEventCount, totalProjectIds, orphanedEventCutoffUtc); } - public async Task DeleteOrphanedEventsByOrganizationAsync(JobContext context) + public Task DeleteOrphanedEventsByOrganizationAsync(JobContext context) + { + return DeleteOrphanedEventsByOrganizationAsync(context, GetOrphanedEventCutoffUtc()); + } + + private async Task DeleteOrphanedEventsByOrganizationAsync(JobContext context, DateTime orphanedEventCutoffUtc) { // get approximate number of unique organization ids var organizationCardinality = await _elasticClient.SearchAsync(s => s .Indices(GetEventIndexPattern()) .Size(0) + .Query(q => RecentEventQuery(q, orphanedEventCutoffUtc)) .AddAggregation("cardinality_organization_id", a => a.Cardinality(c => c.Field(f => f.OrganizationId).PrecisionThreshold(40000)))); double? uniqueOrganizationIdCount = organizationCardinality.Aggregations?.GetCardinality("cardinality_organization_id")?.Value; @@ -201,6 +229,7 @@ public async Task DeleteOrphanedEventsByOrganizationAsync(JobContext context) var organizationIdTerms = await _elasticClient.SearchAsync(s => s .Indices(GetEventIndexPattern()) .Size(0) + .Query(q => RecentEventQuery(q, orphanedEventCutoffUtc)) .AddAggregation("terms_organization_id", a => a.Terms(c => c.Field(f => f.OrganizationId).Include(new TermsInclude(batchNumber, buckets)).Size(batchSize * 2)))); string[] organizationIds = organizationIdTerms.Aggregations?.GetStringTerms("terms_organization_id")?.Buckets.Select(b => b.Key.ToString()!).ToArray() ?? []; @@ -223,10 +252,12 @@ public async Task DeleteOrphanedEventsByOrganizationAsync(JobContext context) _logger.LogInformation("{BatchNumber}/{BatchCount}: Found {OrphanedEventCount} orphaned events from missing organizations {MissingOrganizationIds} out of {OrganizationIdCount}", batchNumber, buckets, missingOrganizationIds.Length, missingOrganizationIds, organizationIds.Length); await _elasticClient.DeleteByQueryAsync(r => r .Indices(GetEventIndexPattern()) - .Query(q => q.Terms(t => t.Field(f => f.OrganizationId).Terms(new TermsQueryField(missingOrganizationIds.Select(FieldValueHelper.ToFieldValue).ToList()))))); + .Query(q => q.Bool(b => b.Filter( + f => f.Terms(t => t.Field(e => e.OrganizationId).Terms(new TermsQueryField(missingOrganizationIds.Select(FieldValueHelper.ToFieldValue).ToList()))), + f => RecentEventQuery(f, orphanedEventCutoffUtc))))); } - _logger.LogInformation("Found {OrphanedEventCount} orphaned events from missing organizations out of {OrganizationIdCount}", totalOrphanedEventCount, totalOrganizationIds); + _logger.LogInformation("Found {OrphanedEventCount} orphaned events from missing organizations out of {OrganizationIdCount} since {OrphanedEventCutoffUtc}", totalOrphanedEventCount, totalOrganizationIds, orphanedEventCutoffUtc); } public async Task FixDuplicateStacks(JobContext context) @@ -401,6 +432,16 @@ private Task RenewLockAsync(JobContext context) return context.RenewLockAsync(); } + private DateTime GetOrphanedEventCutoffUtc() + { + return _timeProvider.GetUtcNow().UtcDateTime.Subtract(OrphanedEventLookback); + } + + private static QueryDescriptor RecentEventQuery(QueryDescriptor query, DateTime orphanedEventCutoffUtc) + { + return query.Range(r => r.Date(d => d.Field(e => e.CreatedUtc).Gte(orphanedEventCutoffUtc))); + } + private string GetEventIndexPattern() { return $"{_config.Events.VersionedName}-*"; diff --git a/tests/Exceptionless.Tests/Jobs/CleanupOrphanedDataJobTests.cs b/tests/Exceptionless.Tests/Jobs/CleanupOrphanedDataJobTests.cs index af3f405b16..f8c6368eb2 100644 --- a/tests/Exceptionless.Tests/Jobs/CleanupOrphanedDataJobTests.cs +++ b/tests/Exceptionless.Tests/Jobs/CleanupOrphanedDataJobTests.cs @@ -6,6 +6,7 @@ using Exceptionless.Tests.Utility; using Foundatio.Repositories; using Foundatio.Repositories.Utility; +using Microsoft.Extensions.Diagnostics.HealthChecks; using Xunit; namespace Exceptionless.Tests.Jobs; @@ -390,6 +391,71 @@ public async Task DeleteOrphanedEventsByOrganization_TwoTenantsOneDeleted_OnlyDe Assert.Equal(120, totalAfter); } + [Fact] + public async Task RunAsync_OrphanedEventsAcrossOccurrenceDates_DeletesOnlyEventsWithinCreatedUtcLookback() + { + var now = DateTimeOffset.UtcNow; + TimeProvider.SetUtcNow(now); + + var organization = await _organizationRepository.AddAsync( + _organizationData.GenerateSampleOrganization(_billingManager, _plans), + o => o.ImmediateConsistency()); + var project = await _projectRepository.AddAsync( + _projectData.GenerateSampleProject(), + o => o.ImmediateConsistency()); + var stack = await _stackRepository.AddAsync( + _stackData.GenerateSampleStack(), + o => o.ImmediateConsistency()); + + var cutoffUtc = TimeProvider.GetUtcNow().UtcDateTime.Subtract(CleanupOrphanedDataJob.OrphanedEventLookback); + var beforeCutoffUtc = cutoffUtc.AddMilliseconds(-1); + var validEvent = _eventData.GenerateEvent(organization.Id, project.Id, stack.Id, occurrenceDate: now); + + string missingStackId = ObjectId.GenerateNewId().ToString(); + var stackOrphanAtCutoff = _eventData.GenerateEvent(organization.Id, project.Id, missingStackId, occurrenceDate: now); + stackOrphanAtCutoff.CreatedUtc = cutoffUtc; + var stackOrphanBeforeCutoff = _eventData.GenerateEvent(organization.Id, project.Id, missingStackId, occurrenceDate: now); + stackOrphanBeforeCutoff.CreatedUtc = beforeCutoffUtc; + var stackOrphanWithOldOccurrence = _eventData.GenerateEvent(organization.Id, project.Id, missingStackId, occurrenceDate: now.Subtract(TimeSpan.FromDays(30))); + stackOrphanWithOldOccurrence.CreatedUtc = now.UtcDateTime; + + string missingProjectId = ObjectId.GenerateNewId().ToString(); + var projectOrphanAtCutoff = _eventData.GenerateEvent(organization.Id, missingProjectId, stack.Id, occurrenceDate: now); + projectOrphanAtCutoff.CreatedUtc = cutoffUtc; + var projectOrphanBeforeCutoff = _eventData.GenerateEvent(organization.Id, missingProjectId, stack.Id, occurrenceDate: now); + projectOrphanBeforeCutoff.CreatedUtc = beforeCutoffUtc; + + string missingOrganizationId = ObjectId.GenerateNewId().ToString(); + var organizationOrphanAtCutoff = _eventData.GenerateEvent(missingOrganizationId, project.Id, stack.Id, occurrenceDate: now); + organizationOrphanAtCutoff.CreatedUtc = cutoffUtc; + var organizationOrphanBeforeCutoff = _eventData.GenerateEvent(missingOrganizationId, project.Id, stack.Id, occurrenceDate: now); + organizationOrphanBeforeCutoff.CreatedUtc = beforeCutoffUtc; + + await _eventRepository.AddAsync([ + validEvent, + stackOrphanAtCutoff, + stackOrphanBeforeCutoff, + stackOrphanWithOldOccurrence, + projectOrphanAtCutoff, + projectOrphanBeforeCutoff, + organizationOrphanAtCutoff, + organizationOrphanBeforeCutoff + ], o => o.ImmediateConsistency()); + + await _job.RunAsync(TestCancellationToken); + + var remainingEvents = await _eventRepository.GetAllAsync(o => o.PageLimit(10).ImmediateConsistency()); + Assert.Equal(4, remainingEvents.Total); + Assert.Contains(remainingEvents.Documents, e => e.Id == validEvent.Id); + Assert.Contains(remainingEvents.Documents, e => e.Id == stackOrphanBeforeCutoff.Id); + Assert.Contains(remainingEvents.Documents, e => e.Id == projectOrphanBeforeCutoff.Id); + Assert.Contains(remainingEvents.Documents, e => e.Id == organizationOrphanBeforeCutoff.Id); + Assert.DoesNotContain(remainingEvents.Documents, e => e.Id == stackOrphanAtCutoff.Id); + Assert.DoesNotContain(remainingEvents.Documents, e => e.Id == stackOrphanWithOldOccurrence.Id); + Assert.DoesNotContain(remainingEvents.Documents, e => e.Id == projectOrphanAtCutoff.Id); + Assert.DoesNotContain(remainingEvents.Documents, e => e.Id == organizationOrphanAtCutoff.Id); + } + [Fact] public async Task FixDuplicateStacks_WithDuplicatesAcrossTenants_MergesCorrectly() { @@ -536,7 +602,7 @@ public async Task RunAsync_NoOrphans_PreservesEverything() } [Fact] - public async Task RunAsync_EmptyDatabase_CompletesWithoutError() + public async Task RunAsync_EmptyDatabase_UpdatesHealth() { // Arrange - nothing @@ -545,5 +611,36 @@ public async Task RunAsync_EmptyDatabase_CompletesWithoutError() var totalAfter = await _eventRepository.CountAsync(o => o.IncludeSoftDeletes().ImmediateConsistency()); Assert.Equal(0, totalAfter); + + var health = await _job.CheckHealthAsync(new HealthCheckContext(), TestCancellationToken); + Assert.Equal(HealthStatus.Healthy, health.Status); + Assert.Equal("Job has run in the last 65 minutes.", health.Description); + } + + [Fact] + public async Task RunAsync_NoRecentEventsOrDuplicateStacks_UpdatesHealth() + { + var now = DateTimeOffset.UtcNow; + TimeProvider.SetUtcNow(now); + + var organization = await _organizationRepository.AddAsync( + _organizationData.GenerateSampleOrganization(_billingManager, _plans), + o => o.ImmediateConsistency()); + var project = await _projectRepository.AddAsync( + _projectData.GenerateProject(organizationId: organization.Id), + o => o.ImmediateConsistency()); + var stack = await _stackRepository.AddAsync( + _stackData.GenerateStack(projectId: project.Id, organizationId: organization.Id), + o => o.ImmediateConsistency()); + var historicalEvent = _eventData.GenerateEvent(organization.Id, project.Id, stack.Id, occurrenceDate: now.Subtract(TimeSpan.FromDays(30))); + historicalEvent.CreatedUtc = now.UtcDateTime.Subtract(CleanupOrphanedDataJob.OrphanedEventLookback).AddMilliseconds(-1); + await _eventRepository.AddAsync(historicalEvent, o => o.ImmediateConsistency()); + + await _job.RunAsync(TestCancellationToken); + + Assert.Equal(1, await _eventRepository.CountAsync(o => o.IncludeSoftDeletes().ImmediateConsistency())); + var health = await _job.CheckHealthAsync(new HealthCheckContext(), TestCancellationToken); + Assert.Equal(HealthStatus.Healthy, health.Status); + Assert.Equal("Job has run in the last 65 minutes.", health.Description); } }