Skip to content
Merged
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
23 changes: 18 additions & 5 deletions tests/Kafdeck.Architecture.Tests/ClusterExplorerServiceTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ public async Task Capability_denial_is_a_limitation_without_changing_healthy_met
[Fact]
public async Task Retryable_refresh_failure_serves_stale_metadata_as_degraded()
{
var clock = new MutableTestTimeProvider(DateTimeOffset.UtcNow);
var calls = 0;
var port = new FakeKafkaAdministrationPort
{
Expand All @@ -72,20 +73,20 @@ public async Task Retryable_refresh_failure_serves_stale_metadata_as_degraded()
var call = Interlocked.Increment(ref calls);
if (call == 1)
{
var oldObservation = Observation(DateTimeOffset.UtcNow - TimeSpan.FromMilliseconds(1500));
var oldObservation = Observation(clock.GetUtcNow() - TimeSpan.FromMilliseconds(1500));
return Task.FromResult(KafkaResult<ClusterMetadata>.Success(Metadata("cluster-a"), oldObservation));
}

return Task.FromResult(KafkaResult<ClusterMetadata>.Failed(
Failure(KafkaFailureCategory.Unavailable, "temporarily_unavailable", true),
Observation()));
Observation(clock.GetUtcNow())));
},
Capabilities = (_, _, _) => Task.FromResult(KafkaResult<KafkaCapabilities>.Success(
new KafkaCapabilities(Array.Empty<KafkaCapabilityStatus>()),
Observation())),
Observation(clock.GetUtcNow()))),
};
var policy = Policy(clusterMetadataTtl: TimeSpan.FromSeconds(1));
var service = new ClusterExplorerService(port, new KafkaSnapshotCoordinator(policy), policy);
var service = new ClusterExplorerService(port, new KafkaSnapshotCoordinator(policy, clock), policy);

var initial = await service.GetClusterAsync("cluster-a");
var stale = await service.GetClusterAsync("cluster-a");
Expand Down Expand Up @@ -149,6 +150,18 @@ private static KafkaFailure Failure(
bool retryable) =>
new(category, code, "Safe failure.", retryable);

private sealed class MutableTestTimeProvider : TimeProvider
{
private readonly DateTimeOffset _utcNow;

public MutableTestTimeProvider(DateTimeOffset utcNow)
{
_utcNow = utcNow;
}

public override DateTimeOffset GetUtcNow() => _utcNow;
}

private sealed class FakeKafkaAdministrationPort : IKafkaAdministrationPort
{
public Func<string, KafkaOperationContext, CancellationToken, Task<KafkaResult<ClusterMetadata>>>? ClusterMetadata { get; init; }
Expand Down Expand Up @@ -196,4 +209,4 @@ public Task<KafkaResult<KafkaCapabilities>> GetCapabilitiesAsync(
new KafkaCapabilities(Array.Empty<KafkaCapabilityStatus>()),
Observation()));
}
}
}
Loading