Skip to content
Open
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
2 changes: 1 addition & 1 deletion src/Exceptionless.Core/Jobs/CleanupDataJob.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,12 +27,15 @@ public RemoveBotEventsWorkItemHandler(IEventRepository eventRepository, ILockPro
public override async Task HandleItemAsync(WorkItemContext context)
{
var wi = context.GetData<RemoveBotEventsWorkItem>()!;
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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,24 +28,33 @@ public ThrottleBotsPlugin(ICacheClient cacheClient, IQueue<WorkItemData> 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<EventContext> 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<int?>(throttleCacheKey, null);
if (requestCount.HasValue)
{
Expand All @@ -62,14 +71,14 @@ public override async Task EventBatchProcessingAsync(ICollection<EventContext> 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)
});
Expand Down
24 changes: 23 additions & 1 deletion src/Exceptionless.Core/Repositories/EventRepository.cs
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ public async Task<bool> UpdateSessionStartLastActivityAsync(string id, DateTime
return true;
}

public Task<long> RemoveAllAsync(string organizationId, string? clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor<PersistentEvent>? options = null)
public Task<long> RemoveAllByOrganizationAndClientIpAsync(string organizationId, string? clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor<PersistentEvent>? options = null)
{
ArgumentException.ThrowIfNullOrEmpty(organizationId);

Expand All @@ -77,6 +77,28 @@ public Task<long> RemoveAllAsync(string organizationId, string? clientIpAddress,
return RemoveAllAsync(q => query, options);
}

public Task<long> RemoveAllByProjectAndClientIpAsync(string organizationId, string projectId, string clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor<PersistentEvent>? options = null)
{
ArgumentException.ThrowIfNullOrWhiteSpace(organizationId);
ArgumentException.ThrowIfNullOrWhiteSpace(projectId);
ArgumentException.ThrowIfNullOrWhiteSpace(clientIpAddress);

var query = new RepositoryQuery<PersistentEvent>()
.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<FindResults<PersistentEvent>> GetByReferenceIdAsync(string projectId, string referenceId)
{
return FindAsync(q => q.Project(projectId).FieldEquals(e => e.ReferenceId, referenceId).SortDescending(e => e.Date), o => o.PageLimit(10));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,8 @@ public interface IEventRepository : IRepositoryOwnedByOrganizationAndProject<Per
Task<PreviousAndNextEventIdResult> GetPreviousAndNextEventIdsAsync(PersistentEvent ev, AppFilter? systemFilter = null, DateTime? utcStart = null, DateTime? utcEnd = null);
Task<FindResults<PersistentEvent>> GetOpenSessionsAsync(DateTime createdBeforeUtc, CommandOptionsDescriptor<PersistentEvent>? options = null);
Task<bool> UpdateSessionStartLastActivityAsync(string id, DateTime lastActivityUtc, bool isSessionEnd = false, bool hasError = false, bool sendNotifications = true);
Task<long> RemoveAllAsync(string organizationId, string? clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor<PersistentEvent>? options = null);
Task<long> RemoveAllByOrganizationAndClientIpAsync(string organizationId, string? clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor<PersistentEvent>? options = null);
Task<long> RemoveAllByProjectAndClientIpAsync(string organizationId, string projectId, string clientIpAddress, DateTime? utcStart, DateTime? utcEnd, CommandOptionsDescriptor<PersistentEvent>? options = null);
Task<long> RemoveAllByStackIdsAsync(string[] stackIds);
}

Expand Down
Original file line number Diff line number Diff line change
@@ -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<IEventRepository, RecordingRepository>();
var handler = new RemoveBotEventsWorkItemHandler(repository, null!, NullLoggerFactory.Instance);
var item = CreateWorkItem();
await handler.HandleItemAsync(CreateContext(item));

var arguments = Assert.IsType<object[]>(((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<IEventRepository, RecordingRepository>();
var handler = new RemoveBotEventsWorkItemHandler(repository, null!, NullLoggerFactory.Instance);
var item = CreateWorkItem() with { ProjectId = projectId! };
await Assert.ThrowsAnyAsync<ArgumentException>(() => 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);
}
}
}
97 changes: 97 additions & 0 deletions tests/Exceptionless.Tests/Plugins/ThrottleBotsPluginTests.cs
Original file line number Diff line number Diff line change
@@ -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<WorkItemData> _queue;
private readonly ThrottleBotsPlugin _plugin;

public ThrottleBotsPluginTests(ITestOutputHelper output) : base(output)
{
_cache = new InMemoryCacheClient(new InMemoryCacheClientOptions { TimeProvider = _time, LoggerFactory = Log });
_queue = new InMemoryQueue<WorkItemData>(new InMemoryQueueOptions<WorkItemData> { TimeProvider = _time, LoggerFactory = Log });
var configuration = new ConfigurationBuilder().AddInMemoryCollection(new Dictionary<string, string?>
{
[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();
}
}
45 changes: 43 additions & 2 deletions tests/Exceptionless.Tests/Repositories/EventRepositoryTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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);
Expand Down