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
81 changes: 49 additions & 32 deletions src/Aspire.Hosting/Dcp/KubernetesService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -88,11 +88,12 @@ internal sealed class KubernetesService(ILogger<KubernetesService> logger, IOpti
private readonly SemaphoreSlim _kubeconfigReadSemaphore = new(1);

private DcpKubernetesClient? _kubernetes;
private int _kubernetesApiReady;
private ResiliencePipeline? _resiliencePipeline;
private bool _disposed;

public TimeSpan MaxRetryDuration { get; set; } = TimeSpan.FromSeconds(20);
public TimeSpan KubernetesConfigInitializationTimeout { get; set; } = TimeSpan.FromSeconds(60);
public TimeSpan KubernetesInitializationTimeout { get; set; } = TimeSpan.FromSeconds(60);

public Task<T> GetAsync<T>(string name, string? namespaceParameter = null, CancellationToken cancellationToken = default)
where T : CustomResource, IKubernetesStaticMetadata
Expand All @@ -102,22 +103,22 @@ public Task<T> GetAsync<T>(string name, string? namespaceParameter = null, Cance
return ExecuteWithRetry(
DcpApiOperationType.Get,
T.ObjectKind,
async (kubernetes) =>
async (kubernetes, operationCancellationToken) =>
{
var response = string.IsNullOrEmpty(namespaceParameter)
? await kubernetes.CustomObjects.GetClusterCustomObjectWithHttpMessagesAsync(
GroupVersion.Group,
GroupVersion.Version,
resourceType,
name,
cancellationToken: cancellationToken).ConfigureAwait(false)
cancellationToken: operationCancellationToken).ConfigureAwait(false)
: await kubernetes.CustomObjects.GetNamespacedCustomObjectWithHttpMessagesAsync(
GroupVersion.Group,
GroupVersion.Version,
namespaceParameter,
resourceType,
name,
cancellationToken: cancellationToken).ConfigureAwait(false);
cancellationToken: operationCancellationToken).ConfigureAwait(false);

return KubernetesJson.Deserialize<T>(response.Body.ToString());
},
Expand All @@ -135,22 +136,22 @@ public Task<T> CreateAsync<T>(T obj, CancellationToken cancellationToken = defau
return ExecuteWithRetry(
DcpApiOperationType.Create,
T.ObjectKind,
async (kubernetes) =>
async (kubernetes, operationCancellationToken) =>
{
var response = string.IsNullOrEmpty(namespaceParameter)
? await kubernetes.CustomObjects.CreateClusterCustomObjectWithHttpMessagesAsync(
obj,
GroupVersion.Group,
GroupVersion.Version,
resourceType,
cancellationToken: cancellationToken).ConfigureAwait(false)
cancellationToken: operationCancellationToken).ConfigureAwait(false)
: await kubernetes.CustomObjects.CreateNamespacedCustomObjectWithHttpMessagesAsync(
obj,
GroupVersion.Group,
GroupVersion.Version,
namespaceParameter,
resourceType,
cancellationToken: cancellationToken).ConfigureAwait(false);
cancellationToken: operationCancellationToken).ConfigureAwait(false);

return KubernetesJson.Deserialize<T>(response.Body.ToString());
},
Expand All @@ -168,7 +169,7 @@ public Task<T> PatchAsync<T>(T obj, V1Patch patch, CancellationToken cancellatio
return ExecuteWithRetry(
DcpApiOperationType.Patch,
T.ObjectKind,
async (kubernetes) =>
async (kubernetes, operationCancellationToken) =>
{
var response = string.IsNullOrEmpty(namespaceParameter)
? await kubernetes.CustomObjects.PatchClusterCustomObjectWithHttpMessagesAsync(
Expand All @@ -177,15 +178,15 @@ public Task<T> PatchAsync<T>(T obj, V1Patch patch, CancellationToken cancellatio
GroupVersion.Version,
resourceType,
obj.Metadata.Name,
cancellationToken: cancellationToken).ConfigureAwait(false)
cancellationToken: operationCancellationToken).ConfigureAwait(false)
: await kubernetes.CustomObjects.PatchNamespacedCustomObjectWithHttpMessagesAsync(
patch,
GroupVersion.Group,
GroupVersion.Version,
namespaceParameter,
resourceType,
obj.Metadata.Name,
cancellationToken: cancellationToken).ConfigureAwait(false);
cancellationToken: operationCancellationToken).ConfigureAwait(false);

return KubernetesJson.Deserialize<T>(response.Body.ToString());
},
Expand All @@ -202,20 +203,20 @@ public Task<List<T>> ListAsync<T>(string? namespaceParameter = null, Cancellatio
return ExecuteWithRetry(
DcpApiOperationType.List,
T.ObjectKind,
async (kubernetes) =>
async (kubernetes, operationCancellationToken) =>
{
var response = string.IsNullOrEmpty(namespaceParameter)
? await kubernetes.CustomObjects.ListClusterCustomObjectWithHttpMessagesAsync(
GroupVersion.Group,
GroupVersion.Version,
resourceType,
cancellationToken: cancellationToken).ConfigureAwait(false)
cancellationToken: operationCancellationToken).ConfigureAwait(false)
: await kubernetes.CustomObjects.ListNamespacedCustomObjectWithHttpMessagesAsync(
GroupVersion.Group,
GroupVersion.Version,
namespaceParameter,
resourceType,
cancellationToken: cancellationToken).ConfigureAwait(false);
cancellationToken: operationCancellationToken).ConfigureAwait(false);

return KubernetesJson.Deserialize<CustomResourceList<T>>(response.Body.ToString()).Items;
},
Expand All @@ -232,22 +233,22 @@ public Task<T> DeleteAsync<T>(string name, string? namespaceParameter = null, Ca
return ExecuteWithRetry(
DcpApiOperationType.Delete,
T.ObjectKind,
async (kubernetes) =>
async (kubernetes, operationCancellationToken) =>
{
var response = string.IsNullOrEmpty(namespaceParameter)
? await kubernetes.CustomObjects.DeleteClusterCustomObjectWithHttpMessagesAsync(
GroupVersion.Group,
GroupVersion.Version,
resourceType,
name,
cancellationToken: cancellationToken).ConfigureAwait(false)
cancellationToken: operationCancellationToken).ConfigureAwait(false)
: await kubernetes.CustomObjects.DeleteNamespacedCustomObjectWithHttpMessagesAsync(
GroupVersion.Group,
GroupVersion.Version,
namespaceParameter,
resourceType,
name,
cancellationToken: cancellationToken).ConfigureAwait(false);
cancellationToken: operationCancellationToken).ConfigureAwait(false);

return KubernetesJson.Deserialize<T>(response.Body.ToString());
},
Expand Down Expand Up @@ -295,7 +296,8 @@ public Task<T> DeleteAsync<T>(string name, string? namespaceParameter = null, Ca
#pragma warning restore CS0618 // Type or member is obsolete
},
RetryOnConnectivityAndConflictErrors,
restartCancellationToken);
restartCancellationToken,
marksApiReady: false);
};

await foreach (var item in PeriodicRestartAsyncEnumerable.CreateAsync(innerWatchFactory, restartInterval: TimeSpan.FromMinutes(5), cancellationToken: cancellationToken).ConfigureAwait(false))
Expand Down Expand Up @@ -343,7 +345,7 @@ public Task<Stream> GetLogStreamAsync<T>(
return ExecuteWithRetry(
DcpApiOperationType.GetLogSubresource,
T.ObjectKind,
async (kubernetes) =>
async (kubernetes, operationCancellationToken) =>
{
var response = await kubernetes.ReadSubResourceAsStreamAsync(
GroupVersion.Group,
Expand All @@ -353,7 +355,7 @@ public Task<Stream> GetLogStreamAsync<T>(
Logs.SubResourceName,
obj.Metadata.Namespace(),
queryParams,
cancellationToken
operationCancellationToken
).ConfigureAwait(false);

return response.Body;
Expand All @@ -370,15 +372,15 @@ public Task StopServerAsync(string resourceCleanup = ResourceCleanup.Full, Cance
return ExecuteWithRetry(
DcpApiOperationType.ServerStop,
DcpExecutionResourceType,
async (kubernetes) =>
async (kubernetes, operationCancellationToken) =>
{
await kubernetes.PatchExecutionDocumentAsync(
new ApiServerExecution
{
ApiServerStatus = ApiServerStatus.Stopping,
ShutdownResourceCleanup = ResourceCleanup.Full
},
cancellationToken
operationCancellationToken
).ConfigureAwait(false);
return (object?)null;
},
Expand All @@ -394,15 +396,15 @@ public async Task CleanupResourcesAsync(CancellationToken cancellationToken = de
var executionDoc = await ExecuteWithRetry(
DcpApiOperationType.ResourceCleanup,
DcpExecutionResourceType,
async (kubernetes) =>
async (kubernetes, operationCancellationToken) =>
{
return await kubernetes.PatchExecutionDocumentAsync(
new ApiServerExecution
{
ApiServerStatus = ApiServerStatus.CleaningResources,
ShutdownResourceCleanup = ResourceCleanup.Full
},
cancellationToken
operationCancellationToken
).ConfigureAwait(false);
},
RetryOnConnectivityErrors,
Expand All @@ -429,7 +431,7 @@ public async Task CleanupResourcesAsync(CancellationToken cancellationToken = de

await retryPipeline.ExecuteAsync(async cancellationContext =>
{
return await _kubernetes!.GetExecutionDocumentAsync(cancellationToken).ConfigureAwait(false);
return await _kubernetes!.GetExecutionDocumentAsync(cancellationContext).ConfigureAwait(false);
}, cancellationToken).ConfigureAwait(false);
}

Expand All @@ -455,25 +457,29 @@ private Task<TResult> ExecuteWithRetry<TResult>(
string resourceType,
Func<DcpKubernetesClient, TResult> operation,
Func<Exception, bool> isRetryable,
CancellationToken cancellationToken)
CancellationToken cancellationToken,
bool marksApiReady = true)
{
return ExecuteWithRetry<TResult>(
operationType,
resourceType,
(DcpKubernetesClient kubernetes) => Task.FromResult(operation(kubernetes)),
(DcpKubernetesClient kubernetes, CancellationToken _) => Task.FromResult(operation(kubernetes)),
isRetryable,
cancellationToken);
cancellationToken,
marksApiReady);
}

private async Task<TResult> ExecuteWithRetry<TResult>(
DcpApiOperationType operationType,
string resourceType,
Func<DcpKubernetesClient, Task<TResult>> operation,
Func<DcpKubernetesClient, CancellationToken, Task<TResult>> operation,
Func<Exception, bool> isRetryable,
CancellationToken cancellationToken)
CancellationToken cancellationToken,
bool marksApiReady = true)
{
using var activity = ProfilingTelemetry.StartDcpKubernetesApi(configuration, operationType, resourceType);
var retryCount = 0;
var useInitializationTimeout = Volatile.Read(ref _kubernetesApiReady) == 0;

try
{
Expand All @@ -491,19 +497,30 @@ async ValueTask EnsureKubernetesClientAsync(CancellationToken cancellationToken)
// The first DCP request must also wait for DCP to create its kubeconfig. Give that initialization
// its own budget so slower hosts do not consume the shorter already-running API retry budget.
var initializationPipeline = CreateKubernetesCallResiliencePipeline(
KubernetesConfigInitializationTimeout,
KubernetesInitializationTimeout,
RetryOnConnectivityErrors,
activity,
() => retryCount++);
await initializationPipeline.ExecuteAsync(EnsureKubernetesClientAsync, cancellationToken).ConfigureAwait(false);
}

var resiliencePipeline = CreateKubernetesCallResiliencePipeline(MaxRetryDuration, isRetryable, activity, () => retryCount++);
// A parsed kubeconfig does not guarantee that DCP can process requests yet. Keep using the startup
// budget until an API operation succeeds, then fail steady-state API calls on the normal budget.
var retryDuration = useInitializationTimeout
? KubernetesInitializationTimeout
: MaxRetryDuration;

var resiliencePipeline = CreateKubernetesCallResiliencePipeline(retryDuration, isRetryable, activity, () => retryCount++);
Comment thread
adamint marked this conversation as resolved.
return await resiliencePipeline.ExecuteAsync(async (cancellationToken) =>
{
// Keep connection establishment inside the retry loop so kubeconfig read failures remain retryable.
await EnsureKubernetesClientAsync(cancellationToken).ConfigureAwait(false);
return await operation(_kubernetes!).ConfigureAwait(false);
var result = await operation(_kubernetes!, cancellationToken).ConfigureAwait(false);
if (marksApiReady)
{
Volatile.Write(ref _kubernetesApiReady, 1);
}
return result;
}, cancellationToken).ConfigureAwait(false);
}
finally
Expand Down
Loading
Loading