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
18 changes: 12 additions & 6 deletions src/Dapr.Actors.Next.Abstractions/Options/DaprActorsOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -52,32 +52,37 @@ public sealed class DaprActorsOptions

/// <summary>
/// Gets or sets the idle timeout for actor activations.
/// <see langword="null"/> leaves the value unset so the Dapr runtime default applies.
/// </summary>
public TimeSpan ActorIdleTimeout { get; set; } = TimeSpan.FromMinutes(60);
public TimeSpan? ActorIdleTimeout { get; set; }

/// <summary>
/// Gets or sets the timeout used when draining ongoing actor calls.
/// <see langword="null"/> leaves the value unset so the Dapr runtime default applies.
/// </summary>
public TimeSpan DrainOngoingCallTimeout { get; private set; } = TimeSpan.FromSeconds(30);
public TimeSpan? DrainOngoingCallTimeout { get; set; }

/// <summary>
/// Gets or sets the timeout used when draining rebalanced actors.
/// <see langword="null"/> leaves the value unset so the Dapr runtime default applies.
/// </summary>
public TimeSpan DrainRebalancedActorsTimeout
public TimeSpan? DrainRebalancedActorsTimeout
{
get => DrainOngoingCallTimeout;
set => DrainOngoingCallTimeout = value;
}

/// <summary>
/// Gets or sets a value indicating whether rebalanced actors drain in-flight calls.
/// <see langword="null"/> leaves the value unset so the Dapr runtime default applies.
/// </summary>
public bool DrainRebalancedActors { get; set; } = true;
public bool? DrainRebalancedActors { get; set; }

/// <summary>
/// Gets or sets a value indicating whether reentrant actor calls are allowed.
/// <see langword="null"/> leaves the value unset so the Dapr runtime default applies.
/// </summary>
public bool EnableReentrancy { get; set; }
public bool? EnableReentrancy { get; set; }

/// <summary>
/// Gets or sets a value indicating whether the sidecar (gRPC) transport is used for actor state,
Expand All @@ -90,8 +95,9 @@ public TimeSpan DrainRebalancedActorsTimeout

/// <summary>
/// Gets or sets the maximum reentrant call depth when reentrancy is enabled.
/// <see langword="null"/> leaves the value unset so the Dapr runtime default applies.
/// </summary>
public int MaxReentrantDepth { get; set; } = 32;
public int? MaxReentrantDepth { get; set; }

internal void CopyFrom(DaprActorsOptions source)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,10 +26,10 @@ public ValidateOptionsResult Validate(string? name, DaprActorsOptions options)
List<string>? failures = null;

AddFailureIf(options.DefaultContractVersion <= 0, "DefaultContractVersion must be greater than zero.", ref failures);
AddFailureIf(options.ActorIdleTimeout <= TimeSpan.Zero, "ActorIdleTimeout must be greater than zero.", ref failures);
AddFailureIf(options.DrainOngoingCallTimeout <= TimeSpan.Zero, "DrainOngoingCallTimeout must be greater than zero.", ref failures);
AddFailureIf(options.DrainRebalancedActorsTimeout <= TimeSpan.Zero, "DrainRebalancedActorsTimeout must be greater than zero.", ref failures);
AddFailureIf(options.MaxReentrantDepth <= 0, "MaxReentrantDepth must be greater than zero.", ref failures);
AddFailureIf(options.ActorIdleTimeout is { } globalIdleTimeout && globalIdleTimeout <= TimeSpan.Zero, "ActorIdleTimeout must be greater than zero.", ref failures);
AddFailureIf(options.DrainOngoingCallTimeout is { } globalDrainTimeout && globalDrainTimeout <= TimeSpan.Zero, "DrainOngoingCallTimeout must be greater than zero.", ref failures);
AddFailureIf(options.DrainRebalancedActorsTimeout is { } globalDrainRebalancedTimeout && globalDrainRebalancedTimeout <= TimeSpan.Zero, "DrainRebalancedActorsTimeout must be greater than zero.", ref failures);
AddFailureIf(options.MaxReentrantDepth is { } globalMaxDepth && globalMaxDepth <= 0, "MaxReentrantDepth must be greater than zero.", ref failures);

foreach (var registration in options.Actors.Registrations)
{
Expand Down
39 changes: 30 additions & 9 deletions src/Dapr.Actors.Next.Core/Transport/DaprActorEventsTransport.cs
Original file line number Diff line number Diff line change
Expand Up @@ -240,17 +240,38 @@ private static P.SubscribeActorEventsRequestInitialAlpha1 CreateInitialRequest(S
return request;
}

var entityConfig = new P.ActorEntityConfig
// Only assign fields the app configured; unset fields keep their
// Dapr runtime defaults.
var entityConfig = new P.ActorEntityConfig();

if (config.ActorIdleTimeout is { } idleTimeout)
{
ActorIdleTimeout = Duration.FromTimeSpan(config.ActorIdleTimeout),
DrainOngoingCallTimeout = Duration.FromTimeSpan(config.DrainOngoingCallTimeout),
DrainRebalancedActors = config.DrainRebalancedActors,
Reentrancy = new P.ActorReentrancyConfig
entityConfig.ActorIdleTimeout = Duration.FromTimeSpan(idleTimeout);
}

if (config.DrainOngoingCallTimeout is { } drainTimeout)
{
entityConfig.DrainOngoingCallTimeout = Duration.FromTimeSpan(drainTimeout);
}

if (config.DrainRebalancedActors is { } drainRebalanced)
{
entityConfig.DrainRebalancedActors = drainRebalanced;
}

if (config.EnableReentrancy is not null || config.MaxReentrantDepth is not null)
{
var reentrancy = new P.ActorReentrancyConfig
{
Enabled = config.EnableReentrancy,
MaxStackDepth = config.MaxReentrantDepth,
},
};
Enabled = config.EnableReentrancy ?? false,
};
if (config.MaxReentrantDepth is { } maxDepth)
{
reentrancy.MaxStackDepth = maxDepth;
}
entityConfig.Reentrancy = reentrancy;
}

entityConfig.Entities.Add(request.Entities);
request.EntitiesConfig.Add(entityConfig);
return request;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,10 +66,12 @@ public static SubscribeActorEventsResponse Failed(string id, string message, uin

/// <summary>
/// Runtime options advertised with the initial actor event stream registration.
/// A <see langword="null"/> value leaves that field unset on the wire so the
/// Dapr runtime default applies.
/// </summary>
public sealed record SubscribeActorEventsInitialConfig(
TimeSpan ActorIdleTimeout,
TimeSpan DrainOngoingCallTimeout,
bool DrainRebalancedActors,
bool EnableReentrancy,
int MaxReentrantDepth);
TimeSpan? ActorIdleTimeout,
TimeSpan? DrainOngoingCallTimeout,
bool? DrainRebalancedActors,
bool? EnableReentrancy,
int? MaxReentrantDepth);
15 changes: 15 additions & 0 deletions test/Dapr.Actors.Next.Abstractions.Test/OptionsTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,21 @@ public void RegisterActor_NullLiteral_ResolvesToNameOverload()
Assert.Null(registration.TypeOptions);
}

[MinimumDaprRuntimeFact("1.18")]
public void Options_DrainRebalancedActorsTimeout_AliasesDrainOngoingCallTimeout()
{
var options = new DaprActorsOptions();
Assert.Null(options.DrainOngoingCallTimeout);
Assert.Null(options.DrainRebalancedActorsTimeout);

options.DrainRebalancedActorsTimeout = TimeSpan.FromSeconds(7);
Assert.Equal(TimeSpan.FromSeconds(7), options.DrainOngoingCallTimeout);
Assert.Equal(TimeSpan.FromSeconds(7), options.DrainRebalancedActorsTimeout);

options.DrainOngoingCallTimeout = null;
Assert.Null(options.DrainRebalancedActorsTimeout);
}

[MinimumDaprRuntimeFact("1.18")]
public void TypeOptions_DrainRebalancedActorsTimeout_AliasesDrainOngoingCallTimeout()
{
Expand Down
22 changes: 11 additions & 11 deletions test/Dapr.Actors.Next.Core.Test/CoreRuntimeTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -117,11 +117,11 @@ public async Task Stream_advertises_global_defaults_in_initial_config()

var config = advertisement.InitialConfig;
Assert.NotNull(config);
Assert.Equal(TimeSpan.FromMinutes(60), config!.ActorIdleTimeout);
Assert.Equal(TimeSpan.FromSeconds(30), config.DrainOngoingCallTimeout);
Assert.True(config.DrainRebalancedActors);
Assert.False(config.EnableReentrancy);
Assert.Equal(32, config.MaxReentrantDepth);
Assert.Null(config!.ActorIdleTimeout);
Assert.Null(config.DrainOngoingCallTimeout);
Assert.Null(config.DrainRebalancedActors);
Assert.Null(config.EnableReentrancy);
Assert.Null(config.MaxReentrantDepth);
}

[MinimumDaprRuntimeFact("1.18")]
Expand Down Expand Up @@ -151,17 +151,17 @@ public async Task Stream_advertises_merged_type_options_per_type()
var overridden = advertisements["Counter"].InitialConfig;
Assert.NotNull(overridden);
Assert.Equal(TimeSpan.FromMinutes(2), overridden!.ActorIdleTimeout);
Assert.Equal(TimeSpan.FromSeconds(30), overridden.DrainOngoingCallTimeout);
Assert.True(overridden.DrainRebalancedActors);
Assert.Null(overridden.DrainOngoingCallTimeout);
Assert.Null(overridden.DrainRebalancedActors);
Assert.True(overridden.EnableReentrancy);
Assert.Equal(8, overridden.MaxReentrantDepth);

var inherited = advertisements["OtherCounter"].InitialConfig;
Assert.NotNull(inherited);
Assert.Equal(TimeSpan.FromMinutes(10), inherited!.ActorIdleTimeout);
Assert.Equal(TimeSpan.FromSeconds(30), inherited.DrainOngoingCallTimeout);
Assert.True(inherited.DrainRebalancedActors);
Assert.False(inherited.EnableReentrancy);
Assert.Null(inherited.DrainOngoingCallTimeout);
Assert.Null(inherited.DrainRebalancedActors);
Assert.Null(inherited.EnableReentrancy);
Assert.Equal(8, inherited.MaxReentrantDepth);
}

Expand Down Expand Up @@ -193,7 +193,7 @@ public async Task Type_options_override_registration_full_options()
Assert.Equal(TimeSpan.FromMinutes(1), merged!.ActorIdleTimeout);
Assert.Equal(TimeSpan.FromSeconds(7), merged.DrainOngoingCallTimeout);
Assert.False(merged.DrainRebalancedActors);
Assert.False(merged.EnableReentrancy);
Assert.Null(merged.EnableReentrancy);
Assert.Equal(3, merged.MaxReentrantDepth);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,32 @@ await stream.WriteAsync(SubscribeActorEventsResponse.RegisteredActors(
Assert.Equal(7, entityConfig.Reentrancy.MaxStackDepth);
}

[MinimumDaprRuntimeFact("1.18")]
public async Task Transport_omits_unset_fields_in_initial_config()
{
var harness = new GeneratedActorEventsHarness();
var transport = new DaprActorEventsTransport(harness.Client);
using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5));

await using var stream = await transport.OpenStreamAsync(cts.Token);
await stream.WriteAsync(SubscribeActorEventsResponse.RegisteredActors(
["Sparse"],
new SubscribeActorEventsInitialConfig(
null,
TimeSpan.FromSeconds(5),
null,
null,
null)), cts.Token);
var initial = await harness.ReceiveAsync(cts.Token);

var entityConfig = Assert.Single(initial.InitialRequest.EntitiesConfig);
Assert.Equal(new[] { "Sparse" }, entityConfig.Entities.ToArray());
Assert.Null(entityConfig.ActorIdleTimeout);
Assert.Equal(5, entityConfig.DrainOngoingCallTimeout.Seconds);
Assert.False(entityConfig.HasDrainRebalancedActors);
Assert.Null(entityConfig.Reentrancy);
}

[MinimumDaprRuntimeFact("1.18")]
public async Task Transport_maps_failures_and_non_invoke_callbacks()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,7 @@ private static void ConfigureActorHost(WebApplicationBuilder builder)
builder.Services.AddDaprActors(options =>
{
options.DrainRebalancedActorsTimeout = TimeSpan.FromSeconds(1);
options.DrainRebalancedActors = true;
options.EnableReentrancy = true;
options.MaxReentrantDepth = 8;
});
Expand Down
Loading