diff --git a/src/Exceptionless.Core/Jobs/CleanupDataJob.cs b/src/Exceptionless.Core/Jobs/CleanupDataJob.cs index 8282a7c6a1..0b043edaeb 100644 --- a/src/Exceptionless.Core/Jobs/CleanupDataJob.cs +++ b/src/Exceptionless.Core/Jobs/CleanupDataJob.cs @@ -579,7 +579,7 @@ private async Task EnforceEventRetentionDaysAsync(Organization organization, int var cutoff = _timeProvider.GetUtcNow().UtcDateTime.Date.SubtractDays(retentionDays); _logger.RetentionEnforcementEventStart(cutoff, organization.Name, organization.Id); - long removedEvents = await _eventRepository.RemoveAllAsync(organization.Id, null, null, cutoff); + long removedEvents = await _eventRepository.RemoveAllByOrganizationAndClientIpAsync(organization.Id, null, null, cutoff); _logger.RetentionEnforcementEventComplete(organization.Name, organization.Id, removedEvents); } diff --git a/src/Exceptionless.Core/Jobs/WorkItemHandlers/RemoveBotEventsWorkItemHandler.cs b/src/Exceptionless.Core/Jobs/WorkItemHandlers/RemoveBotEventsWorkItemHandler.cs index 0e852ee03d..ea4cdf5ee4 100644 --- a/src/Exceptionless.Core/Jobs/WorkItemHandlers/RemoveBotEventsWorkItemHandler.cs +++ b/src/Exceptionless.Core/Jobs/WorkItemHandlers/RemoveBotEventsWorkItemHandler.cs @@ -27,12 +27,15 @@ public RemoveBotEventsWorkItemHandler(IEventRepository eventRepository, ILockPro public override async Task HandleItemAsync(WorkItemContext context) { var wi = context.GetData()!; + ArgumentException.ThrowIfNullOrWhiteSpace(wi.OrganizationId); + ArgumentException.ThrowIfNullOrWhiteSpace(wi.ProjectId); + ArgumentException.ThrowIfNullOrWhiteSpace(wi.ClientIpAddress); using var _ = Log.BeginScope(new ExceptionlessState().Organization(wi.OrganizationId).Project(wi.ProjectId).Tag("Delete").Tag("Bot")); Log.LogInformation("Received remove bot events work item OrganizationId={OrganizationId} ProjectId={ProjectId}, ClientIpAddress={ClientIpAddress}, UtcStartDate={UtcStartDate}, UtcEndDate={UtcEndDate}", wi.OrganizationId, wi.ProjectId, wi.ClientIpAddress, wi.UtcStartDate, wi.UtcEndDate); await context.ReportProgressAsync(0, $"Starting deleting of bot events... OrganizationId={wi.OrganizationId}"); - long deleted = await _eventRepository.RemoveAllAsync(wi.OrganizationId, wi.ClientIpAddress, wi.UtcStartDate, wi.UtcEndDate); + long deleted = await _eventRepository.RemoveAllByProjectAndClientIpAsync(wi.OrganizationId, wi.ProjectId, wi.ClientIpAddress, wi.UtcStartDate, wi.UtcEndDate); await context.ReportProgressAsync(100, $"Bot events deleted: {deleted} OrganizationId={wi.OrganizationId}"); Log.LogInformation("Removed {Deleted} bot events OrganizationId={OrganizationId} ProjectId={ProjectId}, ClientIpAddress={ClientIpAddress}, UtcStartDate={UtcStartDate}, UtcEndDate={UtcEndDate}", deleted, wi.OrganizationId, wi.ProjectId, wi.ClientIpAddress, wi.UtcStartDate, wi.UtcEndDate); } diff --git a/src/Exceptionless.Core/Plugins/EventProcessor/Default/0_ThrottleBotsPlugin.cs b/src/Exceptionless.Core/Plugins/EventProcessor/Default/0_ThrottleBotsPlugin.cs index 3b18c78bcd..fb6c46d3b3 100644 --- a/src/Exceptionless.Core/Plugins/EventProcessor/Default/0_ThrottleBotsPlugin.cs +++ b/src/Exceptionless.Core/Plugins/EventProcessor/Default/0_ThrottleBotsPlugin.cs @@ -28,24 +28,33 @@ public ThrottleBotsPlugin(ICacheClient cacheClient, IQueue workIte _timeProvider = timeProvider; } + private static string CacheKey(string organizationId, string projectId, string clientIpAddress, long period) => + String.Concat("Organization:", organizationId, ":Project:", projectId, ":bot:", period, ":", clientIpAddress); + public override async Task EventBatchProcessingAsync(ICollection contexts) { if (_options.AppMode == AppMode.Development) return; - var firstContext = contexts.First(); - if (!firstContext.Project.DeleteBotDataEnabled || !firstContext.IncludePrivateInformation) - return; - - // Throttle errors by client ip address to no more than X every 5 minutes. - var clientIpAddressGroups = contexts.GroupBy(c => c.Event.GetRequestInfo(_serializer, _logger)?.ClientIpAddress); + // Keep each project's client IP counters and cleanup tasks isolated. + var clientIpAddressGroups = contexts + .Where(c => c.Project.DeleteBotDataEnabled && c.IncludePrivateInformation) + .GroupBy(c => new + { + c.Event.OrganizationId, + c.Event.ProjectId, + ClientIpAddress = c.Event.GetRequestInfo(_serializer, _logger)?.ClientIpAddress + }); foreach (var clientIpAddressGroup in clientIpAddressGroups) { - if (String.IsNullOrEmpty(clientIpAddressGroup.Key) || clientIpAddressGroup.Key.IsPrivateNetwork()) + var scope = clientIpAddressGroup.Key; + if (String.IsNullOrEmpty(scope.ClientIpAddress) || scope.ClientIpAddress.IsPrivateNetwork()) + { continue; + } var clientIpContexts = clientIpAddressGroup.ToList(); - string throttleCacheKey = String.Concat("bot:", clientIpAddressGroup.Key, ":", _timeProvider.GetUtcNow().UtcDateTime.Floor(_throttlingPeriod).Ticks); + string throttleCacheKey = CacheKey(scope.OrganizationId, scope.ProjectId, scope.ClientIpAddress, _timeProvider.GetUtcNow().UtcDateTime.Floor(_throttlingPeriod).Ticks); int? requestCount = await _cache.GetAsync(throttleCacheKey, null); if (requestCount.HasValue) { @@ -62,14 +71,14 @@ public override async Task EventBatchProcessingAsync(ICollection c continue; var utcNow = _timeProvider.GetUtcNow().UtcDateTime; - _logger.LogInformation("Bot throttle triggered. IP: {IP} Time: {ThrottlingPeriod} Project: {ProjectId}", clientIpAddressGroup.Key, utcNow.Floor(_throttlingPeriod), firstContext.Event.ProjectId); + _logger.LogInformation("Bot throttle triggered. IP: {IP} Time: {ThrottlingPeriod} Organization: {OrganizationId} Project: {ProjectId}", scope.ClientIpAddress, utcNow.Floor(_throttlingPeriod), scope.OrganizationId, scope.ProjectId); // The throttle was triggered, go and delete all the errors that triggered the throttle to reduce bot noise in the system await _workItemQueue.EnqueueAsync(new RemoveBotEventsWorkItem { - OrganizationId = firstContext.Event.OrganizationId, - ProjectId = firstContext.Event.ProjectId, - ClientIpAddress = clientIpAddressGroup.Key, + OrganizationId = scope.OrganizationId, + ProjectId = scope.ProjectId, + ClientIpAddress = scope.ClientIpAddress, UtcStartDate = utcNow.Floor(_throttlingPeriod), UtcEndDate = utcNow.Ceiling(_throttlingPeriod) }); diff --git a/src/Exceptionless.Core/Repositories/EventRepository.cs b/src/Exceptionless.Core/Repositories/EventRepository.cs index 3fba37a69d..0a128146d8 100644 --- a/src/Exceptionless.Core/Repositories/EventRepository.cs +++ b/src/Exceptionless.Core/Repositories/EventRepository.cs @@ -59,7 +59,7 @@ public async Task UpdateSessionStartLastActivityAsync(string id, DateTime return true; } - public Task RemoveAllAsync(string organizationId, string? clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor? options = null) + public Task RemoveAllByOrganizationAndClientIpAsync(string organizationId, string? clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor? options = null) { ArgumentException.ThrowIfNullOrEmpty(organizationId); @@ -77,6 +77,28 @@ public Task RemoveAllAsync(string organizationId, string? clientIpAddress, return RemoveAllAsync(q => query, options); } + public Task RemoveAllByProjectAndClientIpAsync(string organizationId, string projectId, string clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor? options = null) + { + ArgumentException.ThrowIfNullOrWhiteSpace(organizationId); + ArgumentException.ThrowIfNullOrWhiteSpace(projectId); + ArgumentException.ThrowIfNullOrWhiteSpace(clientIpAddress); + + var query = new RepositoryQuery() + .Organization(organizationId) + .Project(projectId) + .FieldEquals(EventIndex.Alias.IpAddress, clientIpAddress); + if (utcStart.HasValue || utcEnd.HasValue) + { + query = query.DateRange(utcStart, utcEnd, InferField(e => e.Date)); + } + if (utcStart.HasValue && utcEnd.HasValue) + { + query = query.Index(utcStart, utcEnd); + } + + return RemoveAllAsync(q => query, options); + } + public Task> GetByReferenceIdAsync(string projectId, string referenceId) { return FindAsync(q => q.Project(projectId).FieldEquals(e => e.ReferenceId, referenceId).SortDescending(e => e.Date), o => o.PageLimit(10)); diff --git a/src/Exceptionless.Core/Repositories/Interfaces/IEventRepository.cs b/src/Exceptionless.Core/Repositories/Interfaces/IEventRepository.cs index c7c2272cbc..3d5f891a38 100644 --- a/src/Exceptionless.Core/Repositories/Interfaces/IEventRepository.cs +++ b/src/Exceptionless.Core/Repositories/Interfaces/IEventRepository.cs @@ -11,7 +11,8 @@ public interface IEventRepository : IRepositoryOwnedByOrganizationAndProject GetPreviousAndNextEventIdsAsync(PersistentEvent ev, AppFilter? systemFilter = null, DateTime? utcStart = null, DateTime? utcEnd = null); Task> GetOpenSessionsAsync(DateTime createdBeforeUtc, CommandOptionsDescriptor? options = null); Task UpdateSessionStartLastActivityAsync(string id, DateTime lastActivityUtc, bool isSessionEnd = false, bool hasError = false, bool sendNotifications = true); - Task RemoveAllAsync(string organizationId, string? clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor? options = null); + Task RemoveAllByOrganizationAndClientIpAsync(string organizationId, string? clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor? options = null); + Task RemoveAllByProjectAndClientIpAsync(string organizationId, string projectId, string clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor? options = null); Task RemoveAllByStackIdsAsync(string[] stackIds); } diff --git a/tests/Exceptionless.Tests/Jobs/WorkItemHandlers/RemoveBotEventsWorkItemHandlerTests.cs b/tests/Exceptionless.Tests/Jobs/WorkItemHandlers/RemoveBotEventsWorkItemHandlerTests.cs new file mode 100644 index 0000000000..77fd83b87d --- /dev/null +++ b/tests/Exceptionless.Tests/Jobs/WorkItemHandlers/RemoveBotEventsWorkItemHandlerTests.cs @@ -0,0 +1,72 @@ +using System.Reflection; +using Exceptionless.Core.Jobs.WorkItemHandlers; +using Exceptionless.Core.Models.WorkItems; +using Exceptionless.Core.Repositories; +using Foundatio.Jobs; +using Foundatio.Utility; +using Microsoft.Extensions.Logging.Abstractions; +using Xunit; + +namespace Exceptionless.Tests.Jobs.WorkItemHandlers; + +public sealed class RemoveBotEventsWorkItemHandlerTests +{ + [Fact] + public async Task HandleItemAsync_CleanupTask_RestrictsDeletionToOwningProject() + { + var repository = DispatchProxy.Create(); + var handler = new RemoveBotEventsWorkItemHandler(repository, null!, NullLoggerFactory.Instance); + var item = CreateWorkItem(); + await handler.HandleItemAsync(CreateContext(item)); + + var arguments = Assert.IsType(((RecordingRepository)repository).DeleteArguments); + Assert.Equal(item.OrganizationId, arguments[0]); + Assert.Equal(item.ProjectId, arguments[1]); + Assert.Equal(item.ClientIpAddress, arguments[2]); + Assert.Equal(item.UtcStartDate, arguments[3]); + Assert.Equal(item.UtcEndDate, arguments[4]); + } + + [Theory] + [InlineData(null)] + [InlineData("")] + [InlineData(" ")] + public async Task HandleItemAsync_LegacyTaskWithoutProject_RejectsUnscopedDeletion(string? projectId) + { + var repository = DispatchProxy.Create(); + var handler = new RemoveBotEventsWorkItemHandler(repository, null!, NullLoggerFactory.Instance); + var item = CreateWorkItem() with { ProjectId = projectId! }; + await Assert.ThrowsAnyAsync(() => handler.HandleItemAsync(CreateContext(item))); + Assert.Equal(0, ((RecordingRepository)repository).DeleteCalls); + } + + private static RemoveBotEventsWorkItem CreateWorkItem() => new() + { + OrganizationId = "organization-a", + ProjectId = "project-a", + ClientIpAddress = "203.0.113.10", + UtcStartDate = new DateTime(2026, 9, 21, 2, 25, 0, DateTimeKind.Utc), + UtcEndDate = new DateTime(2026, 9, 21, 2, 30, 0, DateTimeKind.Utc) + }; + + private static WorkItemContext CreateContext(RemoveBotEventsWorkItem item) => + new(item, "test-job", EmptyLock.Empty, TestContext.Current.CancellationToken, static (_, _) => Task.CompletedTask); + + private class RecordingRepository : DispatchProxy + { + public object?[]? DeleteArguments { get; private set; } + public int DeleteCalls { get; private set; } + + protected override object? Invoke(MethodInfo? targetMethod, object?[]? args) + { + if (targetMethod?.Name == nameof(IEventRepository.RemoveAllByProjectAndClientIpAsync)) + { + DeleteCalls++; + DeleteArguments = args; + return Task.FromResult(1L); + } + + throw new NotSupportedException(targetMethod?.Name); + } + } +} diff --git a/tests/Exceptionless.Tests/Plugins/ThrottleBotsPluginTests.cs b/tests/Exceptionless.Tests/Plugins/ThrottleBotsPluginTests.cs new file mode 100644 index 0000000000..7baecf92c3 --- /dev/null +++ b/tests/Exceptionless.Tests/Plugins/ThrottleBotsPluginTests.cs @@ -0,0 +1,97 @@ +using System.Text.Json; +using Exceptionless.Core; +using Exceptionless.Core.Models; +using Exceptionless.Core.Models.Data; +using Exceptionless.Core.Plugins.EventProcessor; +using Exceptionless.Core.Serialization; +using Foundatio.Caching; +using Foundatio.Jobs; +using Foundatio.Queues; +using Foundatio.Serializer; +using Foundatio.Xunit; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.Time.Testing; +using Xunit; + +namespace Exceptionless.Tests.Plugins; + +public sealed class ThrottleBotsPluginTests : TestWithLoggingBase +{ + private readonly FakeTimeProvider _time = new(new DateTimeOffset(2026, 9, 21, 2, 25, 0, TimeSpan.Zero)); + private readonly InMemoryCacheClient _cache; + private readonly InMemoryQueue _queue; + private readonly ThrottleBotsPlugin _plugin; + + public ThrottleBotsPluginTests(ITestOutputHelper output) : base(output) + { + _cache = new InMemoryCacheClient(new InMemoryCacheClientOptions { TimeProvider = _time, LoggerFactory = Log }); + _queue = new InMemoryQueue(new InMemoryQueueOptions { TimeProvider = _time, LoggerFactory = Log }); + var configuration = new ConfigurationBuilder().AddInMemoryCollection(new Dictionary + { + [nameof(AppOptions.BaseURL)] = "http://localhost", + [nameof(AppOptions.AppMode)] = "Production", + [nameof(AppOptions.BotThrottleLimit)] = "2" + }).Build(); + var serializer = new SystemTextJsonSerializer(new JsonSerializerOptions().ConfigureExceptionlessDefaults()); + _plugin = new ThrottleBotsPlugin(_cache, _queue, serializer, _time, + AppOptions.ReadFromConfiguration(configuration), Log); + } + + [Theory] + [InlineData("organization-a", "project-b")] + [InlineData("organization-b", "project-a")] + public async Task EventBatchProcessingAsync_SharedIpAcrossScopes_KeepsCountersSeparate(string organizationId, string projectId) + { + await _plugin.EventBatchProcessingAsync([CreateContext()]); + var other = CreateContext(organizationId, projectId); + await _plugin.EventBatchProcessingAsync([other]); + Assert.False(other.IsDiscarded); + Assert.Equal(0, (await _queue.GetQueueStatsAsync()).Queued); + + var repeated = CreateContext(); + await _plugin.EventBatchProcessingAsync([repeated]); + Assert.True(repeated.IsDiscarded); + Assert.True(repeated.IsCancelled); + Assert.Equal(1, (await _queue.GetQueueStatsAsync()).Queued); + } + + [Fact] + public async Task EventBatchProcessingAsync_MixedProjects_OnlyDiscardsEnabledProjectAtLimit() + { + var disabled = CreateContext(projectId: "disabled", enabled: false); + var first = CreateContext(); + var second = CreateContext(); + var other = CreateContext(projectId: "other"); + await _plugin.EventBatchProcessingAsync([disabled, first, second, other]); + Assert.True(first.IsDiscarded); + Assert.True(second.IsDiscarded); + Assert.False(disabled.IsDiscarded); + Assert.False(other.IsDiscarded); + Assert.Equal(1, (await _queue.GetQueueStatsAsync()).Queued); + } + + [Fact] + public async Task EventBatchProcessingAsync_NextWindow_StartsNewCounter() + { + await _plugin.EventBatchProcessingAsync([CreateContext(), CreateContext()]); + _time.Advance(TimeSpan.FromMinutes(5)); + var next = CreateContext(); + await _plugin.EventBatchProcessingAsync([next]); + Assert.False(next.IsDiscarded); + } + + private static EventContext CreateContext(string organizationId = "organization-a", string projectId = "project-a", bool enabled = true) + { + var ev = new PersistentEvent { Type = Event.KnownTypes.Error }; + ev.AddRequestInfo(new RequestInfo { ClientIpAddress = "203.0.113.10" }); + return new EventContext(ev, new Organization { Id = organizationId }, + new Project { Id = projectId, OrganizationId = organizationId, DeleteBotDataEnabled = enabled }); + } + + public override ValueTask DisposeAsync() + { + _queue.Dispose(); + _cache.Dispose(); + return base.DisposeAsync(); + } +} diff --git a/tests/Exceptionless.Tests/Repositories/EventRepositoryTests.cs b/tests/Exceptionless.Tests/Repositories/EventRepositoryTests.cs index 95c5d2b628..ce5832ae6a 100644 --- a/tests/Exceptionless.Tests/Repositories/EventRepositoryTests.cs +++ b/tests/Exceptionless.Tests/Repositories/EventRepositoryTests.cs @@ -221,7 +221,7 @@ public async Task GetOpenSessionsAsync() } [Fact] - public async Task RemoveAllByClientIpAndDateAsync() + public async Task RemoveAllByOrganizationAndClientIpAsync_WithDateRange_RemovesMatchingEvents() { const string _clientIpAddress = "123.123.12.255"; const int NUMBER_OF_EVENTS_TO_CREATE = 50; @@ -239,12 +239,53 @@ public async Task RemoveAllByClientIpAndDateAsync() Assert.Equal(_clientIpAddress, ri.ClientIpAddress); }); - await _repository.RemoveAllAsync(TestConstants.OrganizationId, _clientIpAddress, DateTime.UtcNow.SubtractDays(3), DateTime.UtcNow.AddDays(2), o => o.ImmediateConsistency()); + await _repository.RemoveAllByOrganizationAndClientIpAsync(TestConstants.OrganizationId, _clientIpAddress, DateTime.UtcNow.SubtractDays(3), DateTime.UtcNow.AddDays(2), o => o.ImmediateConsistency()); events = (await _repository.GetByProjectIdAsync(TestConstants.ProjectId, o => o.PageLimit(NUMBER_OF_EVENTS_TO_CREATE))).Documents.ToList(); Assert.Empty(events); } + [Theory] + [InlineData(true, true)] + [InlineData(true, false)] + [InlineData(false, true)] + [InlineData(false, false)] + public async Task RemoveAllByProjectAndClientIpAsync_WithProjectScope_PreservesOtherScopesAndFilters(bool hasStart, bool hasEnd) + { + const string clientIpAddress = "203.0.113.10"; + var start = DateTime.UtcNow.Date.AddHours(1); + var end = start.AddMinutes(5); + var inside = start.AddMinutes(1); + + PersistentEvent CreateEvent(string organizationId, string projectId, string ip, DateTime date) + { + var ev = _eventData.GenerateEvent(organizationId, projectId, TestConstants.StackId2, + occurrenceDate: date, generateData: false); + ev.AddRequestInfo(new RequestInfo { ClientIpAddress = ip }); + return ev; + } + + var matching = CreateEvent(TestConstants.OrganizationId, TestConstants.ProjectId, clientIpAddress, inside); + var otherProject = CreateEvent(TestConstants.OrganizationId, TestConstants.ProjectIdWithNoRoles, clientIpAddress, inside); + var otherOrganization = CreateEvent(TestConstants.OrganizationId2, TestConstants.ProjectId, clientIpAddress, inside); + var otherIp = CreateEvent(TestConstants.OrganizationId, TestConstants.ProjectId, "203.0.113.11", inside); + var before = CreateEvent(TestConstants.OrganizationId, TestConstants.ProjectId, clientIpAddress, start.AddMinutes(-1)); + var after = CreateEvent(TestConstants.OrganizationId, TestConstants.ProjectId, clientIpAddress, end.AddMinutes(1)); + await _repository.AddAsync([matching, otherProject, otherOrganization, otherIp, before, after], o => o.ImmediateConsistency()); + + long deleted = await _repository.RemoveAllByProjectAndClientIpAsync(TestConstants.OrganizationId, TestConstants.ProjectId, + clientIpAddress, hasStart ? start : null, hasEnd ? end : null, o => o.ImmediateConsistency()); + + Assert.Equal(1 + (hasStart ? 0 : 1) + (hasEnd ? 0 : 1), deleted); + Assert.Null(await _repository.GetByIdAsync(matching.Id)); + foreach (var preserved in new[] { otherProject, otherOrganization, otherIp }) + { + Assert.NotNull(await _repository.GetByIdAsync(preserved.Id)); + } + Assert.Equal(hasStart, await _repository.GetByIdAsync(before.Id) is not null); + Assert.Equal(hasEnd, await _repository.GetByIdAsync(after.Id) is not null); + } + private async Task CreateDataAsync() { var baseDate = DateTime.UtcNow.SubtractHours(1);