diff --git a/src/Directory.Packages.props b/src/Directory.Packages.props
index 1a712572da..8b5fef92b2 100644
--- a/src/Directory.Packages.props
+++ b/src/Directory.Packages.props
@@ -71,7 +71,7 @@
-
+
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/ApprovalFiles/ThroughputCollector_Report_Throughput_Tests.Should_generate_correct_report.approved.txt b/src/Particular.LicensingComponent.UnitTests/ApprovalFiles/ThroughputCollector_Report_Throughput_Tests.Should_generate_correct_report.approved.txt
index ef9a71cd2f..1a2234c3e8 100644
--- a/src/Particular.LicensingComponent.UnitTests/ApprovalFiles/ThroughputCollector_Report_Throughput_Tests.Should_generate_correct_report.approved.txt
+++ b/src/Particular.LicensingComponent.UnitTests/ApprovalFiles/ThroughputCollector_Report_Throughput_Tests.Should_generate_correct_report.approved.txt
@@ -39,12 +39,13 @@
"DateUTC": "2024-04-25",
"MessageCount": 65
}
- ]
+ ],
+ "DailyThroughputFromEndpoint": []
},
{
"QueueName": "Endpoint2",
"NameHash": "936C4C33C71B498D9C60C90F78138AF1ECF3C9F72089F054563E125C94B0A224",
- "Throughput": 65,
+ "Throughput": 64,
"EndpointIndicators": [
"KnownEndpoint"
],
@@ -70,12 +71,13 @@
"MessageCount": 64
}
],
- "DailyThroughputFromMonitoring": []
+ "DailyThroughputFromMonitoring": [],
+ "DailyThroughputFromEndpoint": []
},
{
"QueueName": "Endpoint3",
"NameHash": "8035D54FFAF0523F509245E1556BB1BD18C37B76A2D7DB4C928024B348A13132",
- "Throughput": 57,
+ "Throughput": 47,
"EndpointIndicators": [
"KnownEndpoint"
],
@@ -110,7 +112,8 @@
"DateUTC": "2024-04-25",
"MessageCount": 45
}
- ]
+ ],
+ "DailyThroughputFromEndpoint": []
},
{
"QueueName": "Endpoint4",
@@ -131,7 +134,8 @@
}
],
"DailyThroughputFromAudit": [],
- "DailyThroughputFromMonitoring": []
+ "DailyThroughputFromMonitoring": [],
+ "DailyThroughputFromEndpoint": []
},
{
"QueueName": "Endpoint5",
@@ -143,7 +147,8 @@
"ScopeHash": "",
"DailyThroughputFromBroker": [],
"DailyThroughputFromAudit": [],
- "DailyThroughputFromMonitoring": []
+ "DailyThroughputFromMonitoring": [],
+ "DailyThroughputFromEndpoint": []
}
],
"TotalQueues": 5,
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..9a765b60f6 100644
--- a/src/Particular.LicensingComponent.UnitTests/ThroughputCollector/ThroughputCollector_Report_Throughput_Tests.cs
+++ b/src/Particular.LicensingComponent.UnitTests/ThroughputCollector/ThroughputCollector_Report_Throughput_Tests.cs
@@ -175,8 +175,8 @@ await DataStore.CreateBuilder()
using (Assert.EnterMultipleScope())
{
Assert.That(report.ReportData.Queues.First(w => w.QueueName == "Endpoint1").Throughput, Is.EqualTo(65), $"Incorrect Throughput recorded for Endpoint1");
- Assert.That(report.ReportData.Queues.First(w => w.QueueName == "Endpoint2").Throughput, Is.EqualTo(65), $"Incorrect Throughput recorded for Endpoint2");
- Assert.That(report.ReportData.Queues.First(w => w.QueueName == "Endpoint3").Throughput, Is.EqualTo(57), $"Incorrect Throughput recorded for Endpoint3");
+ Assert.That(report.ReportData.Queues.First(w => w.QueueName == "Endpoint2").Throughput, Is.EqualTo(64), $"Incorrect Throughput recorded for Endpoint2");
+ Assert.That(report.ReportData.Queues.First(w => w.QueueName == "Endpoint3").Throughput, Is.EqualTo(47), $"Incorrect Throughput recorded for Endpoint3");
Assert.That(report.ReportData.TotalQueues, Is.EqualTo(3), $"Incorrect TotalQueues recorded");
}
}
@@ -239,7 +239,7 @@ await DataStore.CreateBuilder()
Assert.That(report.ReportData.Queues[0].QueueName, Is.EqualTo("Endpoint1_"), $"Incorrect Name for Endpoint1");
//even though the names are different, we should have matched on the sanitized name and hence displayed max throughput from the 2 endpoints
- Assert.That(report.ReportData.Queues[0].Throughput, Is.EqualTo(75), $"Incorrect Throughput recorded for Endpoint1");
+ Assert.That(report.ReportData.Queues[0].Throughput, Is.EqualTo(65), $"Incorrect Throughput recorded for Endpoint1");
Assert.That(report.ReportData.TotalQueues, Is.EqualTo(1), $"Incorrect TotalQueues recorded");
}
@@ -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.UnitTests/ThroughputCollector/ThroughputCollector_ThroughputSummary_Tests.cs b/src/Particular.LicensingComponent.UnitTests/ThroughputCollector/ThroughputCollector_ThroughputSummary_Tests.cs
index a8fae5c875..b26e33e70d 100644
--- a/src/Particular.LicensingComponent.UnitTests/ThroughputCollector/ThroughputCollector_ThroughputSummary_Tests.cs
+++ b/src/Particular.LicensingComponent.UnitTests/ThroughputCollector/ThroughputCollector_ThroughputSummary_Tests.cs
@@ -134,8 +134,8 @@ await DataStore.CreateBuilder()
using (Assert.EnterMultipleScope())
{
Assert.That(summary.First(w => w.Name == "Endpoint1").MaxDailyThroughput, Is.EqualTo(65), $"Incorrect MaxDailyThroughput recorded for Endpoint1");
- Assert.That(summary.First(w => w.Name == "Endpoint2").MaxDailyThroughput, Is.EqualTo(65), $"Incorrect MaxDailyThroughput recorded for Endpoint2");
- Assert.That(summary.First(w => w.Name == "Endpoint3").MaxDailyThroughput, Is.EqualTo(57), $"Incorrect MaxDailyThroughput recorded for Endpoint3");
+ Assert.That(summary.First(w => w.Name == "Endpoint2").MaxDailyThroughput, Is.EqualTo(64), $"Incorrect MaxDailyThroughput recorded for Endpoint2");
+ Assert.That(summary.First(w => w.Name == "Endpoint3").MaxDailyThroughput, Is.EqualTo(47), $"Incorrect MaxDailyThroughput recorded for Endpoint3");
}
}
@@ -247,7 +247,7 @@ await DataStore.CreateBuilder()
Assert.That(summary[0].Name, Is.EqualTo("Endpoint1_"), $"Incorrect Name for Endpoint1");
//even though the names are different, we should have matched on the sanitized name and hence displayed max throughput from the 2 endpoints
- Assert.That(summary[0].MaxDailyThroughput, Is.EqualTo(75), $"Incorrect MaxDailyThroughput recorded for Endpoint1");
+ Assert.That(summary[0].MaxDailyThroughput, Is.EqualTo(65), $"Incorrect MaxDailyThroughput recorded for Endpoint1");
}
}
}
\ No newline at end of file
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..91a9841658 100644
--- a/src/Particular.LicensingComponent/Particular.LicensingComponent.csproj
+++ b/src/Particular.LicensingComponent/Particular.LicensingComponent.csproj
@@ -18,6 +18,7 @@
+
diff --git a/src/Particular.LicensingComponent/ThroughputCollector.cs b/src/Particular.LicensingComponent/ThroughputCollector.cs
index e2186f6dda..726b7275d0 100644
--- a/src/Particular.LicensingComponent/ThroughputCollector.cs
+++ b/src/Particular.LicensingComponent/ThroughputCollector.cs
@@ -69,6 +69,8 @@ public async Task> GetThroughputSummary(Cancella
await foreach (var endpointData in GetDistinctEndpointData(null, cancellationToken))
{
+ var dailyThroughput = endpointData.ThroughputData.DailyThroughput();
+
var endpointSummary = new EndpointThroughputSummary
{
Name = endpointData.Name,
@@ -77,9 +79,9 @@ public async Task> GetThroughputSummary(Cancella
ScopeHash = string.IsNullOrEmpty(endpointData.Scope) ? null : OneWayHasher.CalculateOneWayHash(endpointData.Scope),
UserIndicator = endpointData.UserIndicator ?? (endpointData.IsKnownEndpoint ? Contracts.UserIndicator.NServiceBusEndpoint.ToString() : string.Empty),
IsKnownEndpoint = endpointData.IsKnownEndpoint,
- MaxDailyThroughput = endpointData.ThroughputData.MaxDailyThroughput(),
- MonthlyThroughput = endpointData.ThroughputData.MonthlyThroughput(),
- AverageMonthlyThroughput = endpointData.ThroughputData.AverageMonthlyThroughput()
+ MaxDailyThroughput = dailyThroughput.MaxDailyThroughput(),
+ MonthlyThroughput = dailyThroughput.MonthlyThroughput(),
+ AverageMonthlyThroughput = dailyThroughput.AverageMonthlyThroughput()
};
endpointSummaries.Add(endpointSummary);
@@ -144,10 +146,11 @@ public async Task GenerateThroughputReport(string spVersion, DateT
NoDataOrSendOnly = endpointData.ThroughputData.Sum() == 0,
ScopeHash = string.IsNullOrEmpty(endpointData.Scope) ? "" : OneWayHasher.CalculateOneWayHash(endpointData.Scope),
Scope = masker.Mask(endpointData.Scope ?? ""),
- Throughput = endpointData.ThroughputData.MaxDailyThroughput(),
+ Throughput = endpointData.ThroughputData.DailyThroughput().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);
diff --git a/src/Particular.LicensingComponent/ThroughputDataExtensions.cs b/src/Particular.LicensingComponent/ThroughputDataExtensions.cs
index db47a84a0d..84e9ac156a 100644
--- a/src/Particular.LicensingComponent/ThroughputDataExtensions.cs
+++ b/src/Particular.LicensingComponent/ThroughputDataExtensions.cs
@@ -12,44 +12,53 @@ public static IEnumerable FromSource(this List throughputs) => throughputs.SelectMany(t => t).Sum(kvp => kvp.Value);
- public static long MaxDailyThroughput(this List throughputs)
- {
- var items = throughputs.SelectMany(t => t).ToArray();
-
- if (items.Any())
+ public static long MaxDailyThroughput(this Dictionary dailyThroughput)
+ => dailyThroughput switch
{
- return items.Max(kvp => kvp.Value);
- }
-
- return 0;
- }
-
- public static MonthlyThroughput[] MonthlyThroughput(this List throughputs) => [.. throughputs
- .SelectMany(data => data)
- // Older SQL Reports could return a negative value for daily throughput. These are not valid. See https://github.com/Particular/ServiceControl/pull/5404
- .Where(x => x.Value >= 0)
- .GroupBy(x => x.Key, x => x.Value)
- .ToLookup(x => x.Key, x => x.Max())
- .GroupBy(kvp => kvp.Key.ToString("yyyy-MM", CultureInfo.InvariantCulture), x => x.Sum())
- .Select(group => new MonthlyThroughput(group.Key, group.Sum()))];
-
- public static long AverageMonthlyThroughput(this List throughputs)
- {
- if (!throughputs.Any(x => x.Any()))
+ { Count: 0 } => 0,
+ var x => x.Values.Max()
+ };
+
+ public static Dictionary DailyThroughput(this List throughputs) =>
+ throughputs.SelectMany(
+ throughput => throughput.Select(
+ daily => (
+ Source: throughput.ThroughputSource,
+ Date: daily.Key,
+ Throughput: daily.Value
+ )
+ )
+ )
+ // Older SQL Reports could return a negative value for daily throughput. These are not valid. See https://github.com/Particular/ServiceControl/pull/5404
+ .Where(entry => entry.Throughput >= 0)
+ .GroupBy(entry => entry.Date)
+ .ToDictionary(
+ entries => entries.Key,
+ entries => entries.OrderBy(entry => entry.Source switch
+ {
+ ThroughputSource.Endpoint => 0,
+ ThroughputSource.Audit or ThroughputSource.Monitoring => 1,
+ ThroughputSource.Broker => 2,
+ _ => int.MaxValue
+ })
+ .ThenByDescending(entry => entry.Throughput)
+ .Select(entry => entry.Throughput)
+ .First()
+ );
+
+ public static MonthlyThroughput[] MonthlyThroughput(this Dictionary dailyThroughput) => [
+ ..dailyThroughput
+ .GroupBy(kvp => kvp.Key.ToString("yyyy-MM", CultureInfo.InvariantCulture), kvp => kvp.Value)
+ .Select(group => new MonthlyThroughput(group.Key, group.Sum()))
+ ];
+
+
+ public static long AverageMonthlyThroughput(this Dictionary dailyThroughput)
+ => dailyThroughput switch
{
- return 0;
- }
-
- // keep this in sync with the internal licensing calculation
- var maxDailyThroughput = throughputs
- .SelectMany(x => x)
- .Where(x => x.Value >= 0)
- .GroupBy(x => x.Key, x => x.Value)
- .ToLookup(x => x.Key, x => x.Max())
- .ToDictionary(x => x.Key, x => x.Sum());
-
- return (long)Math.Truncate(maxDailyThroughput.Sum(x => x.Value) / (decimal)maxDailyThroughput.Count * 365 / 12);
- }
+ { Count: 0 } => 0,
+ var throughput => (long)Math.Truncate(throughput.Sum(x => x.Value) / (decimal)throughput.Count * 365 / 12)
+ };
public static bool HasDataFromSource(this IDictionary> throughputPerQueue, ThroughputSource source) =>
throughputPerQueue.Any(queueThroughput => queueThroughput.Value.Any(data => data.ThroughputSource == source && data.Count > 0));
diff --git a/src/ServiceControl.AcceptanceTests/Licensing/When_creating_a_usage_report_on_a_broker_transport.cs b/src/ServiceControl.AcceptanceTests/Licensing/When_creating_a_usage_report_on_a_broker_transport.cs
index 2027647c27..fe82971ff0 100644
--- a/src/ServiceControl.AcceptanceTests/Licensing/When_creating_a_usage_report_on_a_broker_transport.cs
+++ b/src/ServiceControl.AcceptanceTests/Licensing/When_creating_a_usage_report_on_a_broker_transport.cs
@@ -63,7 +63,7 @@ await Define()
{
var summary = await this.TryGet>(
"/api/licensing/endpoints",
- items => items.Any(item => item.MaxDailyThroughput == BrokerThroughput));
+ items => items.Any(item => item.Name == SalesQueue));
endpoints = summary.Item;
@@ -112,8 +112,8 @@ await Define()
Assert.That(sales.GetProperty("DailyThroughputFromMonitoring").GetArrayLength(), Is.EqualTo(1),
"Throughput seen by monitoring has to be reported separately from the broker's");
- Assert.That(sales.GetProperty("Throughput").GetInt64(), Is.EqualTo(BrokerThroughput),
- "The reported figure is the highest daily total across the sources");
+ Assert.That(sales.GetProperty("Throughput").GetInt64(), Is.EqualTo(MonitoringThroughput),
+ "The reported figure comes from the preferred source, not the highest across all sources");
}
}
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()
}