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
26 changes: 23 additions & 3 deletions all.sln
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@

Microsoft Visual Studio Solution File, Format Version 12.00
Microsoft Visual Studio Solution File, Format Version 12.00
# Visual Studio Version 17
VisualStudioVersion = 17.3.32929.385
MinimumVisualStudioVersion = 10.0.40219.1
Expand Down Expand Up @@ -273,6 +272,7 @@ EndProject
Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "SecretManagement", "SecretManagement", "{929A8AD2-DB45-B92A-7930-EBDD2DBAF802}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "SecretManagementSample", "examples\SecretManagement\SecretManagementSample\SecretManagementSample.csproj", "{ED74B33F-3CE7-42EB-BAA1-623F43900B15}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Dapr.AI.Microsoft.Extensions.Test", "test\Dapr.AI.Microsoft.Extensions.Test\Dapr.AI.Microsoft.Extensions.Test.csproj", "{86CBB08F-601A-4B0B-87BF-383D391A961C}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Dapr.DistributedLock.Test", "test\Dapr.DistributedLock.Test\Dapr.DistributedLock.Test.csproj", "{2B8E9CAD-F9A2-43B9-BB1C-619CF64476A0}"
Expand All @@ -297,6 +297,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Dapr.Common.Generators.Test
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Dapr.IntegrationTest.DaprClient", "test\Dapr.IntegrationTest.DaprClient\Dapr.IntegrationTest.DaprClient.csproj", "{7DB95C53-19C1-49D2-B0BA-116DAEFC83DB}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Dapr.IntegrationTest.Configuration", "test\Dapr.IntegrationTest.Configuration\Dapr.IntegrationTest.Configuration.csproj", "{193BBE09-E261-4D65-B4CC-18B88DB049D7}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
Expand Down Expand Up @@ -1495,6 +1497,22 @@ Global
{97CAEE0B-4020-4A86-97DA-9900FDF4DFC6}.Release|x64.Build.0 = Release|Any CPU
{97CAEE0B-4020-4A86-97DA-9900FDF4DFC6}.Release|x86.ActiveCfg = Release|Any CPU
{97CAEE0B-4020-4A86-97DA-9900FDF4DFC6}.Release|x86.Build.0 = Release|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Debug|Any CPU.Build.0 = Debug|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Debug|x64.ActiveCfg = Debug|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Debug|x64.Build.0 = Debug|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Debug|x86.ActiveCfg = Debug|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Debug|x86.Build.0 = Debug|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Release|Any CPU.ActiveCfg = Release|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Release|Any CPU.Build.0 = Release|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Release|x64.ActiveCfg = Release|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Release|x64.Build.0 = Release|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Release|x86.ActiveCfg = Release|Any CPU
{193BBE09-E261-4D65-B4CC-18B88DB049D7}.Release|x86.Build.0 = Release|Any CPU
{97CAEE0B-4020-4A86-97DA-9900FDF4DFC6}.Release|x64.ActiveCfg = Release|Any CPU
{97CAEE0B-4020-4A86-97DA-9900FDF4DFC6}.Release|x64.Build.0 = Release|Any CPU
{97CAEE0B-4020-4A86-97DA-9900FDF4DFC6}.Release|x86.ActiveCfg = Release|Any CPU
{97CAEE0B-4020-4A86-97DA-9900FDF4DFC6}.Release|x86.Build.0 = Release|Any CPU
{01A20A89-53A1-4D5B-B563-89E157718474}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{01A20A89-53A1-4D5B-B563-89E157718474}.Debug|Any CPU.Build.0 = Debug|Any CPU
{01A20A89-53A1-4D5B-B563-89E157718474}.Debug|x64.ActiveCfg = Debug|Any CPU
Expand Down Expand Up @@ -1886,7 +1904,6 @@ Global
{F99CE5A1-FDA1-415C-B1E6-C8787734ACD2} = {27C5D71D-0721-4221-9286-B94AB07B58CF}
{B42DD6AA-255C-4606-8A1B-263B26650DED} = {27C5D71D-0721-4221-9286-B94AB07B58CF}
{1BECAC48-1C83-43D2-A91F-9AFA634A68C7} = {27C5D71D-0721-4221-9286-B94AB07B58CF}
{87842296-C78B-42B9-8E96-5054E79CD700} = {DD020B34-460F-455F-8D17-CF4A949F100B}
{929A8AD2-DB45-B92A-7930-EBDD2DBAF802} = {D687DDC4-66C5-4667-9E3A-FD8B78ECAA78}
{ED74B33F-3CE7-42EB-BAA1-623F43900B15} = {929A8AD2-DB45-B92A-7930-EBDD2DBAF802}
{8777EAD2-419B-4683-826A-82B7C1F4F69F} = {8462B106-175A-423A-BA94-BE0D39D0BD8E}
Expand All @@ -1905,6 +1922,9 @@ Global
{DA1F9FE7-6041-4581-B9AD-685034869CE8} = {27C5D71D-0721-4221-9286-B94AB07B58CF}
{E015C5ED-F93F-4DB6-BA01-BC59B7859A60} = {0AF0FE8D-C234-4F04-8514-32206ACE01BD}
{7DB95C53-19C1-49D2-B0BA-116DAEFC83DB} = {8462B106-175A-423A-BA94-BE0D39D0BD8E}
{193BBE09-E261-4D65-B4CC-18B88DB049D7} = {8462B106-175A-423A-BA94-BE0D39D0BD8E}
{87842296-C78B-42B9-8E96-5054E79CD700} = {0AF0FE8D-C234-4F04-8514-32206ACE01BD}
{A3F7C8B2-1D45-4E92-B5F6-3C8D9E0A2F71} = {27C5D71D-0721-4221-9286-B94AB07B58CF}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {65220BF2-EAE1-4CB2-AA58-EBE80768CB40}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,14 +34,16 @@ public static class DaprConfigurationStoreExtension
/// <param name="client">The <see cref="DaprClient"/> used for the request.</param>
/// <param name="sidecarWaitTimeout">The <see cref="TimeSpan"/> used to configure the timeout waiting for Dapr.</param>
/// <param name="metadata">Optional metadata sent to the configuration store.</param>
/// <param name="optional">When true, does not block startup waiting for the sidecar. Configuration is loaded in the background once the sidecar becomes available.</param>
/// <returns>The <see cref="IConfigurationBuilder"/>.</returns>
public static IConfigurationBuilder AddDaprConfigurationStore(
this IConfigurationBuilder configurationBuilder,
string store,
IReadOnlyList<string> keys,
DaprClient client,
TimeSpan sidecarWaitTimeout,
IReadOnlyDictionary<string, string>? metadata = default)
IReadOnlyDictionary<string, string>? metadata = default,
bool optional = false)
{
ArgumentVerifier.ThrowIfNullOrEmpty(store, nameof(store));
ArgumentVerifier.ThrowIfNull(keys, nameof(keys));
Expand All @@ -54,7 +56,8 @@ public static IConfigurationBuilder AddDaprConfigurationStore(
Client = client,
SidecarWaitTimeout = sidecarWaitTimeout,
IsStreaming = false,
Metadata = metadata
Metadata = metadata,
IsOptional = optional
});

return configurationBuilder;
Expand All @@ -71,14 +74,16 @@ public static IConfigurationBuilder AddDaprConfigurationStore(
/// <param name="client">The <see cref="DaprClient"/> used for the request.</param>
/// <param name="sidecarWaitTimeout">The <see cref="TimeSpan"/> used to configure the timeout waiting for Dapr.</param>
/// <param name="metadata">Optional metadata sent to the configuration store.</param>
/// <param name="optional">When true, does not block startup waiting for the sidecar. Configuration is loaded in the background once the sidecar becomes available.</param>
/// <returns>The <see cref="IConfigurationBuilder"/>.</returns>
public static IConfigurationBuilder AddStreamingDaprConfigurationStore(
this IConfigurationBuilder configurationBuilder,
string store,
IReadOnlyList<string> keys,
DaprClient client,
TimeSpan sidecarWaitTimeout,
IReadOnlyDictionary<string, string>? metadata = default)
IReadOnlyDictionary<string, string>? metadata = default,
bool optional = false)
{
ArgumentVerifier.ThrowIfNullOrEmpty(store, nameof(store));
ArgumentVerifier.ThrowIfNull(keys, nameof(keys));
Expand All @@ -91,7 +96,8 @@ public static IConfigurationBuilder AddStreamingDaprConfigurationStore(
Client = client,
SidecarWaitTimeout = sidecarWaitTimeout,
IsStreaming = true,
Metadata = metadata
Metadata = metadata,
IsOptional = optional
});

return configurationBuilder;
Expand Down
149 changes: 136 additions & 13 deletions src/Dapr.Extensions.Configuration/DaprConfigurationStoreProvider.cs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
using System.Threading;
using System.Threading.Tasks;
using Dapr.Client;
using Grpc.Core;
using Microsoft.Extensions.Configuration;

namespace Dapr.Extensions.Configuration;
Expand All @@ -26,14 +27,19 @@ namespace Dapr.Extensions.Configuration;
/// </summary>
internal class DaprConfigurationStoreProvider : ConfigurationProvider, IDisposable
{
private string store;
private IReadOnlyList<string> keys;
private DaprClient daprClient;
private TimeSpan sidecarWaitTimeout;
private bool isStreaming;
private IReadOnlyDictionary<string, string>? metadata;
private CancellationTokenSource cts;
private static readonly TimeSpan DisposeWaitTimeout = TimeSpan.FromSeconds(1);

private readonly string store;
private readonly IReadOnlyList<string> keys;
private readonly DaprClient daprClient;
private readonly TimeSpan sidecarWaitTimeout;
private readonly bool isStreaming;
private readonly bool isOptional;
private readonly IReadOnlyDictionary<string, string>? metadata;
private readonly CancellationTokenSource cts;
private Task loadTask = Task.CompletedTask;
private Task subscribeTask = Task.CompletedTask;
private int disposed;

/// <summary>
/// Constructor.
Expand All @@ -44,30 +50,108 @@ internal class DaprConfigurationStoreProvider : ConfigurationProvider, IDisposab
/// <param name="sidecarWaitTimeout">The <see cref="TimeSpan"/> used to configure the timeout waiting for Dapr.</param>
/// <param name="isStreaming">Determines if the source is streaming or not.</param>
/// <param name="metadata">Optional metadata sent to the configuration store.</param>
/// <param name="isOptional">When true, does not block startup waiting for the sidecar.</param>
public DaprConfigurationStoreProvider(
string store,
IReadOnlyList<string> keys,
DaprClient daprClient,
TimeSpan sidecarWaitTimeout,
bool isStreaming = false,
IReadOnlyDictionary<string, string>? metadata = default)
IReadOnlyDictionary<string, string>? metadata = default,
bool isOptional = false)
{
this.store = store;
this.keys = keys;
this.daprClient = daprClient;
this.sidecarWaitTimeout = sidecarWaitTimeout;
this.isStreaming = isStreaming;
this.isOptional = isOptional;
this.metadata = metadata ?? new Dictionary<string, string>();
this.cts = new CancellationTokenSource();
}

public void Dispose()
{
if (Interlocked.Exchange(ref disposed, 1) != 0)
{
return;
}

cts.Cancel();

var loadTaskCompleted = WaitForBackgroundTask(loadTask);
var subscribeTaskCompleted = WaitForBackgroundTask(subscribeTask);
if (loadTaskCompleted && subscribeTaskCompleted)
{
// Only dispose the CTS after tracked tasks have stopped using cts.Token.
cts.Dispose();
}
}

/// <inheritdoc/>
public override void Load() => LoadAsync().ConfigureAwait(false).GetAwaiter().GetResult();
public override void Load()
{
if (isOptional)
{
Data = new Dictionary<string, string?>(StringComparer.OrdinalIgnoreCase);
loadTask = Task.Run(() => LoadInBackgroundAsync());
}
else
{
LoadAsync().ConfigureAwait(false).GetAwaiter().GetResult();
}
}

private static bool WaitForBackgroundTask(Task task)
{
try
{
return task.Wait(DisposeWaitTimeout);
}
catch
{
// Observe background task exceptions during disposal.
return true;
}
}

private async Task LoadInBackgroundAsync()
{
while (!cts.Token.IsCancellationRequested)
{
try
{
using var tokenSource = new CancellationTokenSource(sidecarWaitTimeout);
using var linked = CancellationTokenSource.CreateLinkedTokenSource(tokenSource.Token, cts.Token);
await daprClient.WaitForSidecarAsync(linked.Token);

await FetchDataAsync();
OnReload();
return;
}
catch (OperationCanceledException) when (cts.Token.IsCancellationRequested)
{
return;
}
catch (OperationCanceledException)
{
// Sidecar wait timed out — retry after delay.
}
catch (DaprException)
{
// Transient Dapr error — retry after delay.
}

try
{
await Task.Delay(sidecarWaitTimeout, cts.Token);
}
catch (OperationCanceledException)
{
return;
}
}
}

private async Task LoadAsync()
{
Expand All @@ -77,6 +161,11 @@ private async Task LoadAsync()
await daprClient.WaitForSidecarAsync(tokenSource.Token);
}

await FetchDataAsync();
}

private async Task FetchDataAsync()
{
if (isStreaming)
{
subscribeTask = Task.Run(async () =>
Expand All @@ -100,12 +189,43 @@ private async Task LoadAsync()
OnReload();
}
}
catch (Exception)
catch (OperationCanceledException) when (cts.Token.IsCancellationRequested)
{
return;
}
catch (RpcException ex) when (cts.Token.IsCancellationRequested && ex.StatusCode == StatusCode.Cancelled)
{
return;
}
catch (Exception ex) when (ex is DaprException or RpcException)
{
// If we catch an exception, try and cancel the subscription so we can connect again.
if (!string.IsNullOrEmpty(id))
{
await daprClient.UnsubscribeConfiguration(store, id);
try
{
await daprClient.UnsubscribeConfiguration(store, id, cts.Token);
}
catch (OperationCanceledException) when (cts.Token.IsCancellationRequested)
{
return;
}
catch (RpcException unsubscribeException) when (cts.Token.IsCancellationRequested && unsubscribeException.StatusCode == StatusCode.Cancelled)
{
return;
}
catch (Exception unsubscribeException) when (unsubscribeException is DaprException or RpcException)
{
// Ignore transient unsubscribe failures and reconnect after the retry delay.
}
}

try
{
await Task.Delay(sidecarWaitTimeout, cts.Token);
}
catch (OperationCanceledException) when (cts.Token.IsCancellationRequested)
{
return;
}
}
}
Expand All @@ -115,10 +235,13 @@ private async Task LoadAsync()
{
// We don't need to worry about ReloadTokens here because it is a constant response.
var getConfigurationResponse = await daprClient.GetConfiguration(store, keys, metadata, cts.Token);
var data = new Dictionary<string, string?>(StringComparer.OrdinalIgnoreCase);
foreach (var item in getConfigurationResponse.Items)
{
Set(item.Key, item.Value.Value);
data[item.Key] = item.Value.Value;
}

Data = data;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -53,9 +53,16 @@ public class DaprConfigurationStoreSource : IConfigurationSource
/// </summary>
public IReadOnlyDictionary<string, string>? Metadata { get; set; } = default;

/// <summary>
/// Gets or sets a value indicating whether this configuration source is optional.
/// When <c>true</c>, the provider will not block startup waiting for the Dapr sidecar and will
/// instead load configuration in the background once the sidecar becomes available.
/// </summary>
public bool IsOptional { get; set; }

/// <inheritdoc/>
public IConfigurationProvider Build(IConfigurationBuilder builder)
{
return new DaprConfigurationStoreProvider(Store, Keys, Client, SidecarWaitTimeout, IsStreaming, Metadata);
return new DaprConfigurationStoreProvider(Store, Keys, Client, SidecarWaitTimeout, IsStreaming, Metadata, IsOptional);
}
}
}
Loading
Loading