From 4c23789f158c0cf9816513a4be28b802eab7a5a2 Mon Sep 17 00:00:00 2001 From: Mike Minutillo Date: Thu, 17 Sep 2026 16:46:41 +0800 Subject: [PATCH 1/2] Store usage data provided by endpoints --- .../ThroughputSource.cs | 3 +- ...ughputCollector_Report_Throughput_Tests.cs | 2 + .../EndpointThroughput/ReportUsage.cs | 49 +++++++++++++++++++ .../Particular.LicensingComponent.csproj | 1 + .../Infrastructure/NServiceBusFactory.cs | 2 + 5 files changed, 56 insertions(+), 1 deletion(-) create mode 100644 src/Particular.LicensingComponent/EndpointThroughput/ReportUsage.cs diff --git a/src/Particular.LicensingComponent.Contracts/ThroughputSource.cs b/src/Particular.LicensingComponent.Contracts/ThroughputSource.cs index cb51165283..85b1319b72 100644 --- a/src/Particular.LicensingComponent.Contracts/ThroughputSource.cs +++ b/src/Particular.LicensingComponent.Contracts/ThroughputSource.cs @@ -4,5 +4,6 @@ public enum ThroughputSource { Broker, Monitoring, - Audit + Audit, + Endpoint } diff --git a/src/Particular.LicensingComponent.UnitTests/ThroughputCollector/ThroughputCollector_Report_Throughput_Tests.cs b/src/Particular.LicensingComponent.UnitTests/ThroughputCollector/ThroughputCollector_Report_Throughput_Tests.cs index a89f3082d0..5264d5b06d 100644 --- a/src/Particular.LicensingComponent.UnitTests/ThroughputCollector/ThroughputCollector_Report_Throughput_Tests.cs +++ b/src/Particular.LicensingComponent.UnitTests/ThroughputCollector/ThroughputCollector_Report_Throughput_Tests.cs @@ -271,6 +271,7 @@ await DataStore.CreateBuilder() [TestCase(ThroughputSource.Audit)] [TestCase(ThroughputSource.Broker)] [TestCase(ThroughputSource.Monitoring)] + [TestCase(ThroughputSource.Endpoint)] public async Task Should_not_include_throughput_after_report_end_date(ThroughputSource source) { // Arrange @@ -293,6 +294,7 @@ await DataStore.CreateBuilder() ThroughputSource.Audit => queue.DailyThroughputFromAudit, ThroughputSource.Broker => queue.DailyThroughputFromBroker, ThroughputSource.Monitoring => queue.DailyThroughputFromMonitoring, + ThroughputSource.Endpoint => queue.DailyThroughputFromEndpoint, _ => throw new ArgumentOutOfRangeException(nameof(source)) }; diff --git a/src/Particular.LicensingComponent/EndpointThroughput/ReportUsage.cs b/src/Particular.LicensingComponent/EndpointThroughput/ReportUsage.cs new file mode 100644 index 0000000000..ec2c331323 --- /dev/null +++ b/src/Particular.LicensingComponent/EndpointThroughput/ReportUsage.cs @@ -0,0 +1,49 @@ +namespace NServiceBus; + +using System.Text.Json.Serialization; +using Particular.LicensingComponent.Contracts; +using Particular.LicensingComponent.Persistence; +using ServiceControl.Transports.BrokerThroughput; + +public class EndpointUsageReport : IMessage +{ + public required string EndpointName { get; set; } + public DateTimeOffset TimeStamp { get; set; } + public long MessagesSuccessfullyProcessed { get; set; } +} + +[Handler] +class EndpointUsageReportHandler( + ILicensingDataStore licensingDataStore, + IBrokerThroughputQuery? brokerThroughputQuery = null +) : IHandleMessages +{ + public async Task Handle(EndpointUsageReport message, IMessageHandlerContext context) + { + var endpointId = new EndpointIdentifier(message.EndpointName, ThroughputSource.Endpoint); + + var endpoint = await licensingDataStore.GetEndpoint(endpointId, context.CancellationToken); + + if (endpoint is null) + { + // TODO: Fill in more of the endpoint details if needed + + endpoint = new Particular.LicensingComponent.Contracts.Endpoint(endpointId) + { + EndpointIndicators = [EndpointIndicator.KnownEndpoint.ToString()], + SanitizedName = brokerThroughputQuery?.SanitizeEndpointName(endpointId.Name) ?? endpointId.Name + }; + + await licensingDataStore.SaveEndpoint(endpoint, context.CancellationToken); + } + + await licensingDataStore.RecordEndpointThroughput( + message.EndpointName, + ThroughputSource.Endpoint, + [new EndpointDailyThroughput(DateOnly.FromDateTime(message.TimeStamp.Date), message.MessagesSuccessfullyProcessed)], + context.CancellationToken); + } +} + +[JsonSerializable(typeof(EndpointUsageReport))] +public partial class UsageReportingSerializationContext : JsonSerializerContext; diff --git a/src/Particular.LicensingComponent/Particular.LicensingComponent.csproj b/src/Particular.LicensingComponent/Particular.LicensingComponent.csproj index d1136120d9..0a1c29082b 100644 --- a/src/Particular.LicensingComponent/Particular.LicensingComponent.csproj +++ b/src/Particular.LicensingComponent/Particular.LicensingComponent.csproj @@ -19,6 +19,7 @@ + diff --git a/src/ServiceControl/Infrastructure/NServiceBusFactory.cs b/src/ServiceControl/Infrastructure/NServiceBusFactory.cs index a9238b6043..fce530c4d8 100644 --- a/src/ServiceControl/Infrastructure/NServiceBusFactory.cs +++ b/src/ServiceControl/Infrastructure/NServiceBusFactory.cs @@ -30,6 +30,7 @@ public static void Configure(Settings.Settings settings, ITransportCustomization } configuration.Handlers.ServiceControlAssembly.AddAll(); + configuration.Handlers.ParticularLicensingComponentAssembly.AddAll(); configuration.EnableFeature(); @@ -66,6 +67,7 @@ public static void Configure(Settings.Settings settings, ITransportCustomization { SagaAuditMessagesSerializationContext.Default, HeartbeatSerializationContext.Default, + UsageReportingSerializationContext.Default, // This is required until we move all known message types over to source generated contexts new DefaultJsonTypeInfoResolver() } From 8238c123bb42c389e88ab36ac12cf4d0218907f6 Mon Sep 17 00:00:00 2001 From: Mike Minutillo Date: Thu, 17 Sep 2026 16:47:03 +0800 Subject: [PATCH 2/2] Include endpoint provided usage data in usage report --- src/Directory.Packages.props | 2 +- src/Particular.LicensingComponent/ThroughputCollector.cs | 3 ++- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/src/Directory.Packages.props b/src/Directory.Packages.props index 1a712572da..adefc762d5 100644 --- a/src/Directory.Packages.props +++ b/src/Directory.Packages.props @@ -71,7 +71,7 @@ - + diff --git a/src/Particular.LicensingComponent/ThroughputCollector.cs b/src/Particular.LicensingComponent/ThroughputCollector.cs index e0b718c7f2..7cd388cc15 100644 --- a/src/Particular.LicensingComponent/ThroughputCollector.cs +++ b/src/Particular.LicensingComponent/ThroughputCollector.cs @@ -146,7 +146,8 @@ public async Task GenerateThroughputReport(string spVersion, DateT Throughput = endpointData.ThroughputData.MaxDailyThroughput(), DailyThroughputFromAudit = endpointData.ThroughputData.FromSource(ThroughputSource.Audit).Select(s => new DailyThroughput { DateUTC = s.DateUTC, MessageCount = s.MessageCount }).ToArray(), DailyThroughputFromMonitoring = endpointData.ThroughputData.FromSource(ThroughputSource.Monitoring).Select(s => new DailyThroughput { DateUTC = s.DateUTC, MessageCount = s.MessageCount }).ToArray(), - DailyThroughputFromBroker = notAnNsbEndpoint ? [] : endpointData.ThroughputData.FromSource(ThroughputSource.Broker).Select(s => new DailyThroughput { DateUTC = s.DateUTC, MessageCount = s.MessageCount }).ToArray() + DailyThroughputFromBroker = notAnNsbEndpoint ? [] : endpointData.ThroughputData.FromSource(ThroughputSource.Broker).Select(s => new DailyThroughput { DateUTC = s.DateUTC, MessageCount = s.MessageCount }).ToArray(), + DailyThroughputFromEndpoint = endpointData.ThroughputData.FromSource(ThroughputSource.Endpoint).Select(s => new DailyThroughput { DateUTC = s.DateUTC, MessageCount = s.MessageCount }).ToArray() }; queueThroughputs.Add(queueThroughput);