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/Directory.Packages.props
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,7 @@
<PackageVersion Include="OpenTelemetry.Instrumentation.Http" Version="1.18.0" />
<PackageVersion Include="OpenTelemetry.Instrumentation.Runtime" Version="1.18.0" />
<PackageVersion Include="Particular.Approvals" Version="2.0.1" />
<PackageVersion Include="Particular.LicensingComponent.Report" Version="1.3.0" />
<PackageVersion Include="Particular.LicensingComponent.Report" Version="1.4.0" />
<PackageVersion Include="Particular.Licensing.Sources" Version="7.4.0" />
<PackageVersion Include="Particular.Obsoletes" Version="1.1.0" />
<PackageVersion Include="Particular.ServicePulse.Core" Version="2.12.0" />
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,5 +4,6 @@ public enum ThroughputSource
{
Broker,
Monitoring,
Audit
Audit,
Endpoint
}
Original file line number Diff line number Diff line change
Expand Up @@ -39,12 +39,13 @@
"DateUTC": "2024-04-25",
"MessageCount": 65
}
]
],
"DailyThroughputFromEndpoint": []
},
{
"QueueName": "Endpoint2",
"NameHash": "936C4C33C71B498D9C60C90F78138AF1ECF3C9F72089F054563E125C94B0A224",
"Throughput": 65,
"Throughput": 64,
"EndpointIndicators": [
"KnownEndpoint"
],
Expand All @@ -70,12 +71,13 @@
"MessageCount": 64
}
],
"DailyThroughputFromMonitoring": []
"DailyThroughputFromMonitoring": [],
"DailyThroughputFromEndpoint": []
},
{
"QueueName": "Endpoint3",
"NameHash": "8035D54FFAF0523F509245E1556BB1BD18C37B76A2D7DB4C928024B348A13132",
"Throughput": 57,
"Throughput": 47,
"EndpointIndicators": [
"KnownEndpoint"
],
Expand Down Expand Up @@ -110,7 +112,8 @@
"DateUTC": "2024-04-25",
"MessageCount": 45
}
]
],
"DailyThroughputFromEndpoint": []
},
{
"QueueName": "Endpoint4",
Expand All @@ -131,7 +134,8 @@
}
],
"DailyThroughputFromAudit": [],
"DailyThroughputFromMonitoring": []
"DailyThroughputFromMonitoring": [],
"DailyThroughputFromEndpoint": []
},
{
"QueueName": "Endpoint5",
Expand All @@ -143,7 +147,8 @@
"ScopeHash": "",
"DailyThroughputFromBroker": [],
"DailyThroughputFromAudit": [],
"DailyThroughputFromMonitoring": []
"DailyThroughputFromMonitoring": [],
"DailyThroughputFromEndpoint": []
}
],
"TotalQueues": 5,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
}
Expand Down Expand Up @@ -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");
}
Expand Down Expand Up @@ -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
Expand All @@ -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))
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
}

Expand Down Expand Up @@ -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");
}
}
}
Original file line number Diff line number Diff line change
@@ -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<EndpointUsageReport>
{
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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
// TODO: Fill in more of the endpoint details if needed
// TODO: Fill in more of the endpoint details if needed, e.g. Scope if it becomes available


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;
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
</ItemGroup>

<ItemGroup>
<PackageReference Include="NServiceBus" />
<PackageReference Include="Particular.LicensingComponent.Report" />
</ItemGroup>

Expand Down
13 changes: 8 additions & 5 deletions src/Particular.LicensingComponent/ThroughputCollector.cs
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,8 @@ public async Task<List<EndpointThroughputSummary>> GetThroughputSummary(Cancella

await foreach (var endpointData in GetDistinctEndpointData(null, cancellationToken))
{
var dailyThroughput = endpointData.ThroughputData.DailyThroughput();

var endpointSummary = new EndpointThroughputSummary
{
Name = endpointData.Name,
Expand All @@ -77,9 +79,9 @@ public async Task<List<EndpointThroughputSummary>> 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);
Expand Down Expand Up @@ -144,10 +146,11 @@ public async Task<SignedReport> 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);
Expand Down
81 changes: 45 additions & 36 deletions src/Particular.LicensingComponent/ThroughputDataExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -12,44 +12,53 @@ public static IEnumerable<EndpointDailyThroughput> FromSource(this List<Throughp

public static long Sum(this List<ThroughputData> throughputs) => throughputs.SelectMany(t => t).Sum(kvp => kvp.Value);

public static long MaxDailyThroughput(this List<ThroughputData> throughputs)
{
var items = throughputs.SelectMany(t => t).ToArray();

if (items.Any())
public static long MaxDailyThroughput(this Dictionary<DateOnly, long> dailyThroughput)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

alias Dictionary<DateOnly, long> to DailyThroughputValues or something similar

=> dailyThroughput switch
{
return items.Max(kvp => kvp.Value);
}

return 0;
}

public static MonthlyThroughput[] MonthlyThroughput(this List<ThroughputData> 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<ThroughputData> throughputs)
{
if (!throughputs.Any(x => x.Any()))
{ Count: 0 } => 0,
var x => x.Values.Max()
};

public static Dictionary<DateOnly, long> DailyThroughput(this List<ThroughputData> 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<DateOnly, long> 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<DateOnly, long> dailyThroughput)
=> dailyThroughput switch
{
return 0;
}

// keep this in sync with the internal licensing calculation

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this comment was important, but the calculation is spread out now. I think it still needs to be captured somewhere that we have an internal calculation that this class needs to match behaviour-wise

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<string, IEnumerable<ThroughputData>> throughputPerQueue, ThroughputSource source) =>
throughputPerQueue.Any(queueThroughput => queueThroughput.Value.Any(data => data.ThroughputSource == source && data.Count > 0));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ await Define<Context>()
{
var summary = await this.TryGet<List<EndpointThroughputSummary>>(
"/api/licensing/endpoints",
items => items.Any(item => item.MaxDailyThroughput == BrokerThroughput));
items => items.Any(item => item.Name == SalesQueue));

endpoints = summary.Item;

Expand Down Expand Up @@ -112,8 +112,8 @@ await Define<Context>()
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");
}
}

Expand Down
2 changes: 2 additions & 0 deletions src/ServiceControl/Infrastructure/NServiceBusFactory.cs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ public static void Configure(Settings.Settings settings, ITransportCustomization
}

configuration.Handlers.ServiceControlAssembly.AddAll();
configuration.Handlers.ParticularLicensingComponentAssembly.AddAll();

configuration.EnableFeature<RegisterPluginMessagesFeature>();

Expand Down Expand Up @@ -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()
}
Expand Down
Loading