Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@

namespace Kafdeck.Infrastructure.Kafka;

public sealed class ConfluentKafkaConsumerGroupReadAdapter : IConsumerGroupReadPort, IDisposable
public sealed class ConfluentKafkaConsumerGroupReadAdapter : IConsumerGroupReadPort, IConsumerGroupSamplingReadPort, IDisposable
{
private readonly KafkaAdminClientRegistry _clients;
private readonly TimeProvider _timeProvider;
Expand Down Expand Up @@ -68,6 +68,110 @@ public Task<ReadViewResult<IReadOnlyList<ConsumerGroupSummary>>> ListGroupsAsync
.ToArray();
});

public Task<ReadViewResult<ConsumerGroupPage>> ListGroupPageAsync(
string clusterId,
string? afterGroupId,
int maxItems,
ReadViewOperationContext operation,
CancellationToken cancellationToken)
{
if (maxItems is < 1 or > 2_000)
{
throw new ArgumentOutOfRangeException(
nameof(maxItems));
}

if (afterGroupId is not null &&
(string.IsNullOrWhiteSpace(afterGroupId) ||
!string.Equals(
afterGroupId,
afterGroupId.Trim(),
StringComparison.Ordinal) ||
afterGroupId.Any(char.IsControl)))
{
throw new ArgumentException(
"Consumer-group sampling cursor is invalid.",
nameof(afterGroupId));
}

return ExecuteAsync(
clusterId,
operation,
cancellationToken,
async (client, timeout, token) =>
{
var result =
await client.ListConsumerGroupsAsync(
new ListConsumerGroupsOptions
{
RequestTimeout = timeout,
})
.WaitAsync(token)
.ConfigureAwait(false);

var selected =
new SortedSet<string>(
StringComparer.Ordinal);
foreach (var group in result.Valid)
{
var groupId =
group.GroupId;
if (afterGroupId is not null &&
string.CompareOrdinal(
groupId,
afterGroupId) <= 0)
{
continue;
}

selected.Add(
groupId);
if (selected.Count >
maxItems)
{
selected.Remove(
selected.Max!);
}
}

var items =
selected
.Take(maxItems)
.ToArray();

EnsureStringBudget(
items,
operation.MaxResponseBytes);

string? nextCursor = null;
if (items.Length ==
maxItems)
{
var last =
items[^1];
if (result.Valid.Any(group =>
string.CompareOrdinal(
group.GroupId,
last) > 0))
{
nextCursor =
last;
}
}

return new ConsumerGroupPage(
items
.Select(groupId =>
new ConsumerGroupSummary(
groupId,
CoreConsumerGroupState.Unknown,
null,
false))
.ToArray(),
nextCursor);
});
}

public Task<ReadViewResult<ConsumerGroupDetail>> GetGroupAsync(
string clusterId,
string groupId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,12 @@
namespace Kafdeck.Infrastructure.Persistence;

public sealed class AdoHistoricalMetricMaintenanceStore :
IHistoricalMetricMaintenanceStore
IHistoricalMetricMaintenanceStore,
IHistoricalMetricSamplingLeaseStore
{
private const int SchemaVersion = 2;
private const int SingletonId = 1;
private const int SamplingSingletonId = 2;
private const string Component =
"historical-metrics-maintenance";
private const string RollupSource =
Expand Down Expand Up @@ -460,9 +462,11 @@ EXECUTE FUNCTION kafdeck_enforce_history_exact_writer_fence()
.ConfigureAwait(false);
}

await using (var seed =
connection.CreateCommand())
foreach (var singletonId in
new[] { SingletonId, SamplingSingletonId })
{
await using var seed =
connection.CreateCommand();
seed.Transaction = transaction;
seed.CommandText =
"""
Expand All @@ -484,7 +488,7 @@ DO NOTHING
AddParameter(
seed,
"@singleton_id",
SingletonId);
singletonId);
AddParameter(
seed,
"@updated_at_utc",
Expand Down Expand Up @@ -556,12 +560,39 @@ await transaction
.ConfigureAwait(false);
}

public async Task<HistoricalMetricMaintenanceLease?>
public Task<HistoricalMetricMaintenanceLease?>
TryAcquireLeaseAsync(
string ownerId,
DateTimeOffset nowUtc,
TimeSpan leaseDuration,
CancellationToken cancellationToken = default)
CancellationToken cancellationToken = default) =>
TryAcquireLeaseAsync(
SingletonId,
ownerId,
nowUtc,
leaseDuration,
cancellationToken);

public Task<HistoricalMetricMaintenanceLease?>
TryAcquireSamplingLeaseAsync(
string ownerId,
DateTimeOffset nowUtc,
TimeSpan leaseDuration,
CancellationToken cancellationToken = default) =>
TryAcquireLeaseAsync(
SamplingSingletonId,
ownerId,
nowUtc,
leaseDuration,
cancellationToken);

private async Task<HistoricalMetricMaintenanceLease?>
TryAcquireLeaseAsync(
int singletonId,
string ownerId,
DateTimeOffset nowUtc,
TimeSpan leaseDuration,
CancellationToken cancellationToken)
{
ValidateOwner(ownerId);

Expand Down Expand Up @@ -595,6 +626,7 @@ await connection
await LockAndReadLeaseAsync(
connection,
transaction,
singletonId,
cancellationToken)
.ConfigureAwait(false);

Expand Down Expand Up @@ -649,7 +681,7 @@ UPDATE kafdeck_historical_metric_maintenance
AddParameter(
update,
"@singleton_id",
SingletonId);
singletonId);

if (await update
.ExecuteNonQueryAsync(cancellationToken)
Expand Down Expand Up @@ -961,9 +993,20 @@ private static CancellationTokenRegistration
sqliteConnection);
}

private Task<LeaseRow> LockAndReadLeaseAsync(
DbConnection connection,
DbTransaction transaction,
CancellationToken cancellationToken) =>
LockAndReadLeaseAsync(
connection,
transaction,
SingletonId,
cancellationToken);

private async Task<LeaseRow> LockAndReadLeaseAsync(
DbConnection connection,
DbTransaction transaction,
int singletonId,
CancellationToken cancellationToken)
{
if (!_connectionFactory.SupportsSelectForUpdate)
Expand All @@ -980,7 +1023,7 @@ UPDATE kafdeck_historical_metric_maintenance
AddParameter(
lockCommand,
"@singleton_id",
SingletonId);
singletonId);
await lockCommand
.ExecuteNonQueryAsync(cancellationToken)
.ConfigureAwait(false);
Expand Down Expand Up @@ -1011,7 +1054,7 @@ FROM kafdeck_historical_metric_maintenance
AddParameter(
command,
"@singleton_id",
SingletonId);
singletonId);

await using var reader =
await command
Expand Down
Loading
Loading