From c21f035dde7a9b6dd222555b960deacccf680c38 Mon Sep 17 00:00:00 2001 From: Whit Waldo Date: Thu, 24 Sep 2026 04:01:09 -0500 Subject: [PATCH] Adding a test to prove that reentrant invocation on the same activation works in Actors.Next Signed-off-by: Whit Waldo --- .../RealSidecarActorTests.cs | 113 +++++++++++++++++- 1 file changed, 111 insertions(+), 2 deletions(-) diff --git a/test/Dapr.IntegrationTest.Actors.Next/RealSidecarActorTests.cs b/test/Dapr.IntegrationTest.Actors.Next/RealSidecarActorTests.cs index cdadfbe44..675d79240 100644 --- a/test/Dapr.IntegrationTest.Actors.Next/RealSidecarActorTests.cs +++ b/test/Dapr.IntegrationTest.Actors.Next/RealSidecarActorTests.cs @@ -343,6 +343,31 @@ public async Task Explicit_state_save_persists_before_turn_completion() Assert.Equal(3, final.SchemaVersion); } + /// + /// Verifies a reminder observes state written by a reentrant invocation on the same activation. + /// + [MinimumDaprRuntimeFact("1.18")] + public async Task Reentrant_state_write_is_observed_by_reminder() + { + using var cts = new CancellationTokenSource(Timeout); + var actorId = $"reentrant-state-{Guid.NewGuid():N}"; + + await client!.InvokeActorAsync(CreateInvoke("ProtocolActor", actorId, "StartStateReminder", ReadOnlyMemory.Empty), cancellationToken: cts.Token); + await WaitUntilAsync( + async () => await ReadReminderStateAsync(actorId, cts.Token) >= 1, + cts.Token, + () => "The initial reminder callback did not persist state."); + + await client.InvokeActorAsync(CreateInvoke("ProtocolActor", actorId, "WriteStateReentrantly", Encoding.UTF8.GetBytes("10")), cancellationToken: cts.Token); + + await WaitUntilAsync( + async () => await ReadReminderStateAsync(actorId, cts.Token) >= 11, + cts.Token, + () => "A reminder did not observe the value persisted by the reentrant invocation."); + + await client.InvokeActorAsync(CreateInvoke("ProtocolActor", actorId, "StopStateReminder", ReadOnlyMemory.Empty), cancellationToken: cts.Token); + } + private static P.InvokeActorRequest CreateInvoke( string actorType, string actorId, @@ -391,6 +416,14 @@ private async Task WaitForActorRuntimeAsync() } } + private async Task ReadReminderStateAsync(string actorId, CancellationToken cancellationToken) + { + var response = await client!.InvokeActorAsync( + CreateInvoke("ProtocolActor", actorId, "ReadReminderState", ReadOnlyMemory.Empty), + cancellationToken: cancellationToken); + return JsonSerializer.Deserialize(response.Data.Span); + } + private static async Task WaitUntilAsync(Func predicate, CancellationToken cancellationToken, Func failure) { while (!predicate()) @@ -403,6 +436,19 @@ private static async Task WaitUntilAsync(Func predicate, CancellationToken await Task.Delay(TimeSpan.FromMilliseconds(100), CancellationToken.None); } } + + private static async Task WaitUntilAsync(Func> predicate, CancellationToken cancellationToken, Func failure) + { + while (!await predicate()) + { + if (cancellationToken.IsCancellationRequested) + { + throw new TimeoutException(failure()); + } + + await Task.Delay(TimeSpan.FromMilliseconds(100), CancellationToken.None); + } + } } /// @@ -745,10 +791,11 @@ public Task TimerFiredAsync(CancellationToken cancellationToken) /// /// Handles a real Dapr reminder callback. /// - public Task ReminderAsync(CancellationToken cancellationToken) + public async Task ReminderAsync(CancellationToken cancellationToken) { Probe.Reminders.Add(Id.Value); - return Task.CompletedTask; + var state = await State.GetOrCreateAsync("reentrant-state", static () => 0, cancellationToken); + state.Value++; } /// @@ -782,6 +829,62 @@ public async Task SaveProfileThenReadRawAsync(string name, Cance return envelope!.Value; } + /// + /// Initializes state and schedules callbacks that read and increment it. + /// + public async Task StartStateReminderAsync(CancellationToken cancellationToken) + { + await State.SetAsync("reentrant-state", 0, cancellationToken); + await reminderScheduler.ScheduleAsync( + "ProtocolActor", + Id, + "state-reentrancy", + TimeSpan.FromMilliseconds(100), + TimeSpan.FromMilliseconds(100), + JsonSerializer.Serialize("state-reentrancy"), + overwrite: true, + cancellationToken: cancellationToken); + } + + /// + /// Invokes a same-actor method through the runtime's reentrancy path. + /// + public async Task WriteStateReentrantlyAsync(int value, CancellationToken cancellationToken) + { + await runtime.InvokeAsync( + "ProtocolActor", + Id.Value, + "SetReminderState", + JsonSerializer.SerializeToUtf8Bytes(value), + new Dictionary(), + cancellationToken); + } + + /// + /// Sets the state read by reminder callbacks. + /// + public async Task SetReminderStateAsync(int value, CancellationToken cancellationToken) + { + await State.SetAsync("reentrant-state", value, cancellationToken); + } + + /// + /// Reads the state maintained by reminder callbacks. + /// + public async Task ReadReminderStateAsync(CancellationToken cancellationToken) + { + var state = await State.GetOrCreateAsync("reentrant-state", static () => 0, cancellationToken); + return state.Value; + } + + /// + /// Stops the reminder used by the reentrancy state test. + /// + public async Task StopStateReminderAsync(CancellationToken cancellationToken) + { + await reminderScheduler.CancelAsync("ProtocolActor", Id, "state-reentrancy", cancellationToken); + } + /// /// Forces deactivation through the runtime. /// @@ -812,9 +915,15 @@ public async ValueTask DispatchAsync(IActor actor, ActorD "ManageReminders" => new ActorDispatchResponse(JsonSerializer.SerializeToUtf8Bytes(await protocol.ManageRemindersAsync(cancellationToken))), "TimerFired" => await CompleteAsync(protocol.TimerFiredAsync(cancellationToken)), "reminder" => await CompleteAsync(protocol.ReminderAsync(cancellationToken)), + "state-reentrancy" => await CompleteAsync(protocol.ReminderAsync(cancellationToken)), "SeedLegacy" => await CompleteAsync(protocol.SeedLegacyAsync(JsonSerializer.Deserialize(request.Payload.Span)!, cancellationToken)), "ReadCurrent" => new ActorDispatchResponse(JsonSerializer.SerializeToUtf8Bytes(await protocol.ReadCurrentAsync(cancellationToken))), "SaveProfileThenReadRaw" => new ActorDispatchResponse(JsonSerializer.SerializeToUtf8Bytes(await protocol.SaveProfileThenReadRawAsync(JsonSerializer.Deserialize(request.Payload.Span)!, cancellationToken))), + "StartStateReminder" => await CompleteAsync(protocol.StartStateReminderAsync(cancellationToken)), + "WriteStateReentrantly" => await CompleteAsync(protocol.WriteStateReentrantlyAsync(JsonSerializer.Deserialize(request.Payload.Span), cancellationToken)), + "SetReminderState" => await CompleteAsync(protocol.SetReminderStateAsync(JsonSerializer.Deserialize(request.Payload.Span), cancellationToken)), + "ReadReminderState" => new ActorDispatchResponse(JsonSerializer.SerializeToUtf8Bytes(await protocol.ReadReminderStateAsync(cancellationToken))), + "StopStateReminder" => await CompleteAsync(protocol.StopStateReminderAsync(cancellationToken)), "Deactivate" => await CompleteAsync(protocol.ForceDeactivateAsync(cancellationToken)), _ => throw new InvalidOperationException($"Unknown method '{request.MethodName}'."), };