using Centron.BusinessLogic.Telemetry; using Centron.Data.Entities.Telemetry; using Centron.Host.AspNetCore.Telemetry; using Centron.Interfaces.BL; using JetBrains.Annotations; using NSubstitute; namespace Centron.Tests.BL.Telemetry; [TestSubject(typeof(TelemetryAggregator))] public class TelemetryAggregatorTest { private static readonly DateTime _bucket = new(2026, 5, 2, 12, 15, 0, DateTimeKind.Utc); private static (TelemetryAggregator Aggregator, TelemetryBL BlMock) CreateAggregator() { var bl = Substitute.For((Centron.DAO.DAOSession?)null); bl.UpsertMcpToolUsageBatch(Arg.Any>()) .Returns(Result.AsSuccess()); bl.UpsertApiCallBatch(Arg.Any>()) .Returns(Result.AsSuccess()); // Lookup-resolution always returns one I3D per distinct name (auto-numbered, deterministic). var toolNameCounter = 0; bl.ResolveMcpToolNameI3Ds(Arg.Any>()) .Returns(call => { var names = call.Arg>(); var dict = new Dictionary(StringComparer.Ordinal); foreach (var n in names.Distinct(StringComparer.Ordinal)) dict[n] = ++toolNameCounter; return Result>.AsSuccess(dict); }); var methodNameCounter = 0; bl.ResolveApiMethodNameI3Ds(Arg.Any>()) .Returns(call => { var names = call.Arg>(); var dict = new Dictionary(StringComparer.Ordinal); foreach (var n in names.Distinct(StringComparer.Ordinal)) dict[n] = ++methodNameCounter; return Result>.AsSuccess(dict); }); var hardwareIdCounter = 0; bl.ResolveHardwareIDI3Ds(Arg.Any>()) .Returns(call => { var names = call.Arg>(); var dict = new Dictionary(StringComparer.Ordinal); foreach (var n in names.Distinct(StringComparer.Ordinal)) dict[n] = ++hardwareIdCounter; return Result>.AsSuccess(dict); }); bl.LoadAllHardwareIDs() .Returns(Result>.AsSuccess(new List())); var aggregator = new TelemetryAggregator(action => action(bl)); return (aggregator, bl); } #region SnapToBucket [Theory] [InlineData(0, 0)] [InlineData(7, 0)] [InlineData(14, 0)] [InlineData(15, 15)] [InlineData(29, 15)] [InlineData(30, 30)] [InlineData(44, 30)] [InlineData(45, 45)] [InlineData(59, 45)] public void SnapToBucket_RoundsDownTo15MinuteBoundary(int inputMinute, int expectedMinute) { var input = new DateTime(2026, 5, 2, 12, inputMinute, 0, DateTimeKind.Utc); var snapped = TelemetryAggregator.SnapToBucket(input); Assert.Equal(expectedMinute, snapped.Minute); Assert.Equal(12, snapped.Hour); } [Fact] public void SnapToBucket_PreservesUtcKind() { var input = new DateTime(2026, 5, 2, 12, 17, 0, DateTimeKind.Utc); var snapped = TelemetryAggregator.SnapToBucket(input); Assert.Equal(DateTimeKind.Utc, snapped.Kind); } [Fact] public void SnapToBucket_StripsSecondsAndMilliseconds() { var input = new DateTime(2026, 5, 2, 12, 14, 59, 999, DateTimeKind.Utc); var snapped = TelemetryAggregator.SnapToBucket(input); Assert.Equal(0, snapped.Second); Assert.Equal(0, snapped.Millisecond); Assert.Equal(0, snapped.Minute); } #endregion #region IncrementToolUsage [Fact] public async Task IncrementToolUsage_MergesIdenticalKeysIntoSingleBucket() { var (aggregator, bl) = CreateAggregator(); IReadOnlyCollection? captured = null; bl.UpsertMcpToolUsageBatch(Arg.Do>(b => captured = b)) .Returns(Result.AsSuccess()); for (int i = 0; i < 5; i++) aggregator.IncrementToolUsage(42, "create_ticket", McpToolMode.User, _bucket); await aggregator.FlushAsync(CancellationToken.None); Assert.NotNull(captured); var item = Assert.Single(captured); Assert.Equal(42, item.UserID); Assert.True(item.ToolNameI3D > 0); Assert.Equal(McpToolMode.User, item.ToolMode); Assert.Equal(_bucket, item.BucketStartUtc); Assert.Equal(5, item.Increment); } [Fact] public async Task IncrementToolUsage_DistinctKeys_CreateSeparateBuckets() { var (aggregator, bl) = CreateAggregator(); IReadOnlyCollection? captured = null; bl.UpsertMcpToolUsageBatch(Arg.Do>(b => captured = b)) .Returns(Result.AsSuccess()); aggregator.IncrementToolUsage(1, "tool", McpToolMode.User, _bucket); aggregator.IncrementToolUsage(2, "tool", McpToolMode.User, _bucket); aggregator.IncrementToolUsage(3, "tool", McpToolMode.User, _bucket); await aggregator.FlushAsync(CancellationToken.None); Assert.NotNull(captured); Assert.Equal(3, captured.Count); Assert.All(captured, c => Assert.Equal(1, c.Increment)); Assert.Equal(new[] { 1, 2, 3 }, captured.Select(c => c.UserID).OrderBy(x => x)); } [Fact] public async Task IncrementToolUsage_DifferentBucketStarts_CreateSeparateBuckets() { var (aggregator, bl) = CreateAggregator(); IReadOnlyCollection? captured = null; bl.UpsertMcpToolUsageBatch(Arg.Do>(b => captured = b)) .Returns(Result.AsSuccess()); var laterBucket = _bucket.AddMinutes(15); aggregator.IncrementToolUsage(7, "tool", McpToolMode.User, _bucket); aggregator.IncrementToolUsage(7, "tool", McpToolMode.User, laterBucket); await aggregator.FlushAsync(CancellationToken.None); Assert.NotNull(captured); Assert.Equal(2, captured.Count); Assert.All(captured, c => Assert.Equal(1, c.Increment)); } [Theory] [InlineData(0, "tool")] [InlineData(-1, "tool")] [InlineData(1, "")] public async Task IncrementToolUsage_IgnoresInvalidInput(int userId, string toolName) { var (aggregator, bl) = CreateAggregator(); aggregator.IncrementToolUsage(userId, toolName, McpToolMode.User, _bucket); await aggregator.FlushAsync(CancellationToken.None); bl.DidNotReceive().UpsertMcpToolUsageBatch(Arg.Any>()); } [Fact] public async Task IncrementToolUsage_IgnoresNullToolName() { var (aggregator, bl) = CreateAggregator(); aggregator.IncrementToolUsage(1, null!, McpToolMode.User, _bucket); await aggregator.FlushAsync(CancellationToken.None); bl.DidNotReceive().UpsertMcpToolUsageBatch(Arg.Any>()); } [Fact] public async Task IncrementToolUsage_IsThreadSafe_UnderConcurrentLoad() { var (aggregator, bl) = CreateAggregator(); IReadOnlyCollection? captured = null; bl.UpsertMcpToolUsageBatch(Arg.Do>(b => captured = b)) .Returns(Result.AsSuccess()); Parallel.For(0, 10_000, _ => aggregator.IncrementToolUsage(99, "tool", McpToolMode.User, _bucket)); await aggregator.FlushAsync(CancellationToken.None); Assert.NotNull(captured); var item = Assert.Single(captured); Assert.Equal(10_000, item.Increment); } #endregion #region IncrementApiCall [Theory] [InlineData(0, "method")] [InlineData(-1, "method")] [InlineData(1, "")] public async Task IncrementApiCall_IgnoresInvalidInput(int userId, string methodName) { var (aggregator, bl) = CreateAggregator(); aggregator.IncrementApiCall(userId, TelemetryUserKind.User, TelemetryLicenseKind.Centron, methodName, _bucket); await aggregator.FlushAsync(CancellationToken.None); bl.DidNotReceive().UpsertApiCallBatch(Arg.Any>()); } [Fact] public async Task IncrementApiCall_IgnoresNullMethodName() { var (aggregator, bl) = CreateAggregator(); aggregator.IncrementApiCall(1, TelemetryUserKind.User, TelemetryLicenseKind.Centron, null!, _bucket); await aggregator.FlushAsync(CancellationToken.None); bl.DidNotReceive().UpsertApiCallBatch(Arg.Any>()); } [Fact] public async Task IncrementApiCall_NullLicenseKind_AggregatedSeparatelyFromKnownLicense() { var (aggregator, bl) = CreateAggregator(); IReadOnlyCollection? captured = null; bl.UpsertApiCallBatch(Arg.Do>(b => captured = b)) .Returns(Result.AsSuccess()); // Both calls use the SAME UserKind so the test isolates the LicenseKind splitter // (would still pass even if UserKind were the discriminator if we varied it). aggregator.IncrementApiCall(1, TelemetryUserKind.User, null, "method", _bucket); aggregator.IncrementApiCall(1, TelemetryUserKind.User, TelemetryLicenseKind.Centron, "method", _bucket); await aggregator.FlushAsync(CancellationToken.None); Assert.NotNull(captured); Assert.Equal(2, captured.Count); Assert.Contains(captured, c => c.LicenseKind == null); Assert.Contains(captured, c => c.LicenseKind == TelemetryLicenseKind.Centron); } #endregion #region FlushAsync [Fact] public async Task FlushAsync_OnEmptyBuckets_DoesNotCallBl() { var (aggregator, bl) = CreateAggregator(); await aggregator.FlushAsync(CancellationToken.None); bl.DidNotReceive().UpsertMcpToolUsageBatch(Arg.Any>()); bl.DidNotReceive().UpsertApiCallBatch(Arg.Any>()); } [Fact] public async Task FlushAsync_PassesAggregatedIncrementsToBl() { var (aggregator, bl) = CreateAggregator(); IReadOnlyCollection? tools = null; IReadOnlyCollection? apis = null; bl.UpsertMcpToolUsageBatch(Arg.Do>(b => tools = b)) .Returns(Result.AsSuccess()); bl.UpsertApiCallBatch(Arg.Do>(b => apis = b)) .Returns(Result.AsSuccess()); aggregator.IncrementToolUsage(1, "a", McpToolMode.User, _bucket); aggregator.IncrementToolUsage(2, "b", McpToolMode.User, _bucket); aggregator.IncrementToolUsage(3, "c", McpToolMode.User, _bucket); aggregator.IncrementApiCall(1, TelemetryUserKind.User, TelemetryLicenseKind.Centron, "m1", _bucket); aggregator.IncrementApiCall(2, TelemetryUserKind.WebAccount, TelemetryLicenseKind.Centron, "m2", _bucket); await aggregator.FlushAsync(CancellationToken.None); Assert.NotNull(tools); Assert.NotNull(apis); Assert.Equal(3, tools.Count); Assert.Equal(2, apis.Count); Assert.Contains(tools, t => t.UserID == 2); Assert.Contains(apis, a => a.UserKind == TelemetryUserKind.WebAccount); } [Fact] public async Task FlushAsync_BatchesAt200ItemsPerCall() { var (aggregator, bl) = CreateAggregator(); var batchSizes = new List(); bl.UpsertMcpToolUsageBatch(Arg.Do>(b => batchSizes.Add(b.Count))) .Returns(Result.AsSuccess()); for (int i = 0; i < 450; i++) aggregator.IncrementToolUsage(i + 1, "tool", McpToolMode.User, _bucket); await aggregator.FlushAsync(CancellationToken.None); Assert.Equal(new[] { 200, 200, 50 }, batchSizes); } [Fact] public async Task FlushAsync_OnBlError_RestoresSnapshotIntoBuckets() { var (aggregator, bl) = CreateAggregator(); var responses = new Queue(); responses.Enqueue(Result.AsError("boom")); responses.Enqueue(Result.AsSuccess()); IReadOnlyCollection? secondCallCaptured = null; bl.UpsertMcpToolUsageBatch(Arg.Any>()) .Returns(call => { var arg = call.Arg>(); var result = responses.Dequeue(); if (result.Status == ResultStatus.Success) secondCallCaptured = arg; return result; }); aggregator.IncrementToolUsage(1, "tool", McpToolMode.User, _bucket); aggregator.IncrementToolUsage(1, "tool", McpToolMode.User, _bucket); await aggregator.FlushAsync(CancellationToken.None); await aggregator.FlushAsync(CancellationToken.None); Assert.NotNull(secondCallCaptured); var item = Assert.Single(secondCallCaptured); Assert.Equal(2, item.Increment); } [Fact] public async Task FlushAsync_OnBlError_RestoreCombinesWithIncrementsArrivedDuringFlush() { var (aggregator, bl) = CreateAggregator(); var blocker = new ManualResetEventSlim(false); var midFlight = new ManualResetEventSlim(false); var firstCall = true; bl.UpsertMcpToolUsageBatch(Arg.Any>()) .Returns(_ => { if (firstCall) { firstCall = false; midFlight.Set(); blocker.Wait(); return Result.AsError("boom"); } return Result.AsSuccess(); }); aggregator.IncrementToolUsage(1, "tool", McpToolMode.User, _bucket); var flushTask = Task.Run(() => aggregator.FlushAsync(CancellationToken.None)); midFlight.Wait(); aggregator.IncrementToolUsage(1, "tool", McpToolMode.User, _bucket); blocker.Set(); await flushTask; IReadOnlyCollection? captured = null; bl.UpsertMcpToolUsageBatch(Arg.Do>(b => captured = b)) .Returns(Result.AsSuccess()); await aggregator.FlushAsync(CancellationToken.None); Assert.NotNull(captured); var item = Assert.Single(captured); Assert.Equal(2, item.Increment); } #endregion }