diff --git a/src/Dapr.Workflow.Abstractions/WorkflowRuntimeOptions.cs b/src/Dapr.Workflow.Abstractions/WorkflowRuntimeOptions.cs index 7476953fc..af198ec41 100644 --- a/src/Dapr.Workflow.Abstractions/WorkflowRuntimeOptions.cs +++ b/src/Dapr.Workflow.Abstractions/WorkflowRuntimeOptions.cs @@ -1,4 +1,4 @@ -// ------------------------------------------------------------------------ +// ------------------------------------------------------------------------ // Copyright 2022 The Dapr Authors // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -25,6 +25,9 @@ public sealed class WorkflowRuntimeOptions private readonly List> _registrationActions = []; private int _maxConcurrentWorkflows = 100; private int _maxConcurrentActivities = 100; + private TimeSpan? _historyCacheTtl; + private int _historyCacheMaxInstances; + private long _historyCacheMaxBytes; /// /// Gets or sets the gRPC channel options used for connecting to the Dapr sidecar. @@ -45,14 +48,35 @@ public sealed class WorkflowRuntimeOptions /// when it has gone idle (no turn) for longer than this. null (the default) uses the built-in /// default of one hour. Ignored when is true. /// - public TimeSpan? HistoryCacheTtl { get; set; } + public TimeSpan? HistoryCacheTtl + { + get => _historyCacheTtl; + set + { + if (value is { } configured && configured <= TimeSpan.Zero) + { + throw new ArgumentOutOfRangeException(nameof(value), configured, + "The history cache TTL must be positive."); + } + + _historyCacheTtl = value; + } + } /// /// Gets or sets the maximum number of per-instance histories retained on a single work-item stream; /// least-recently-used entries are evicted beyond it. 0 (the default) uses the built-in default. /// Ignored when is true. /// - public int HistoryCacheMaxInstances { get; set; } + public int HistoryCacheMaxInstances + { + get => _historyCacheMaxInstances; + set + { + ArgumentOutOfRangeException.ThrowIfNegative(value); + _historyCacheMaxInstances = value; + } + } /// /// Gets or sets the byte budget for cached histories on a single work-item stream; least-recently-used @@ -60,7 +84,15 @@ public sealed class WorkflowRuntimeOptions /// and ). Ignored when /// is true. /// - public long HistoryCacheMaxBytes { get; set; } + public long HistoryCacheMaxBytes + { + get => _historyCacheMaxBytes; + set + { + ArgumentOutOfRangeException.ThrowIfNegative(value); + _historyCacheMaxBytes = value; + } + } /// /// Gets the maximum number of concurrent workflow instances that can be executed at the same time. diff --git a/src/Dapr.Workflow/Worker/Grpc/WorkflowHistoryCache.cs b/src/Dapr.Workflow/Worker/Grpc/WorkflowHistoryCache.cs index e03310eb7..ed02aac6f 100644 --- a/src/Dapr.Workflow/Worker/Grpc/WorkflowHistoryCache.cs +++ b/src/Dapr.Workflow/Worker/Grpc/WorkflowHistoryCache.cs @@ -1,4 +1,4 @@ -// ------------------------------------------------------------------------ +// ------------------------------------------------------------------------ // Copyright 2026 The Dapr Authors // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -38,11 +38,13 @@ private sealed class Entry /// Matches the width of the running total and the configured budget, so a large /// history cannot overflow and corrupt the accounting. public required long Bytes { get; init; } + public required LinkedListNode RecencyNode { get; init; } public DateTime LastAccess { get; set; } } private readonly object _lock = new(); private readonly Dictionary _entries = new(); + private readonly LinkedList _recency = new(); private readonly TimeSpan _ttl; private readonly int _maxInstances; private readonly long _maxBytes; @@ -67,16 +69,24 @@ public int Generation } } - /// Initializes the cache. Non-positive ttl/maxInstances use defaults; maxBytes <= 0 means unlimited. + /// Initializes the cache. Null ttl and zero maxInstances use defaults; zero maxBytes means unlimited. public WorkflowHistoryCache( TimeSpan? ttl = null, int maxInstances = 0, long maxBytes = 0, Func? clock = null) { - _ttl = ttl is { } configured && configured > TimeSpan.Zero ? configured : DefaultTtl; - _maxInstances = maxInstances > 0 ? maxInstances : DefaultMaxInstances; - _maxBytes = maxBytes > 0 ? maxBytes : 0; + if (ttl is { } configured && configured <= TimeSpan.Zero) + { + throw new ArgumentOutOfRangeException(nameof(ttl), configured, "The history cache TTL must be positive."); + } + + ArgumentOutOfRangeException.ThrowIfNegative(maxInstances); + ArgumentOutOfRangeException.ThrowIfNegative(maxBytes); + + _ttl = ttl ?? DefaultTtl; + _maxInstances = maxInstances == 0 ? DefaultMaxInstances : maxInstances; + _maxBytes = maxBytes; _clock = clock ?? (() => DateTime.UtcNow); } @@ -90,7 +100,7 @@ public WorkflowHistoryCache( return null; } - entry.LastAccess = _clock(); + TouchLocked(entry); return entry.Events; } } @@ -125,12 +135,20 @@ public void Put(string instanceId, IEnumerable events, int generat if (_entries.TryGetValue(instanceId, out var existing)) { - _totalBytes -= existing.Bytes; + _entries.Remove(instanceId); + RemoveEntryLocked(existing); } - _entries[instanceId] = new Entry { Events = snapshot, Bytes = bytes, LastAccess = _clock() }; + var recencyNode = _recency.AddFirst(instanceId); + _entries[instanceId] = new Entry + { + Events = snapshot, + Bytes = bytes, + LastAccess = _clock(), + RecencyNode = recencyNode + }; _totalBytes += bytes; - EvictToFit(instanceId); + EvictToFit(); } } @@ -161,6 +179,7 @@ public void Reset() lock (_lock) { _entries.Clear(); + _recency.Clear(); _totalBytes = 0; _generation++; } @@ -210,20 +229,33 @@ internal long TotalBytes } } + private void TouchLocked(Entry entry) + { + entry.LastAccess = _clock(); + _recency.Remove(entry.RecencyNode); + _recency.AddFirst(entry.RecencyNode); + } + private void RemoveLocked(string instanceId) { if (_entries.Remove(instanceId, out var entry)) { - _totalBytes -= entry.Bytes; + RemoveEntryLocked(entry); } } + private void RemoveEntryLocked(Entry entry) + { + _recency.Remove(entry.RecencyNode); + _totalBytes -= entry.Bytes; + } + /// /// Evicts least-recently-used entries until within the count and byte bounds. Always keeps the /// just-touched entry so the active working set is never evicted; a lone entry over the byte /// budget is kept (a soft overage) rather than thrashing. /// - private void EvictToFit(string keep) + private void EvictToFit() { while (_entries.Count > 1) { @@ -234,34 +266,13 @@ private void EvictToFit(string keep) return; } - var victim = LeastRecentlyUsedExcept(keep); + var victim = _recency.Last; if (victim is null) { return; } - RemoveLocked(victim); + RemoveLocked(victim.Value); } } - - private string? LeastRecentlyUsedExcept(string keep) - { - string? oldestId = null; - var oldestAccess = DateTime.MaxValue; - foreach (var (instanceId, entry) in _entries) - { - if (instanceId == keep) - { - continue; - } - - if (oldestId is null || entry.LastAccess < oldestAccess) - { - oldestId = instanceId; - oldestAccess = entry.LastAccess; - } - } - - return oldestId; - } } diff --git a/test/Dapr.IntegrationTest.Workflow/StatefulHistory/RequiresDaprHeadFactAttribute.cs b/test/Dapr.IntegrationTest.Workflow/StatefulHistory/RequiresDaprHeadFactAttribute.cs deleted file mode 100644 index eb99fbac2..000000000 --- a/test/Dapr.IntegrationTest.Workflow/StatefulHistory/RequiresDaprHeadFactAttribute.cs +++ /dev/null @@ -1,58 +0,0 @@ -// ------------------------------------------------------------------------ -// Copyright 2026 The Dapr Authors -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// http://www.apache.org/licenses/LICENSE-2.0 -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. -// ------------------------------------------------------------------------ - -using System.Runtime.CompilerServices; - -namespace Dapr.IntegrationTest.Workflow.StatefulHistory; - -/// -/// Marks a test that needs a daprd built from dapr master rather than a release, and skips it -/// otherwise. -/// -/// -/// cannot -/// express this: the feature under test is on no release at all, and that attribute treats an -/// unset or non-semver DAPR_RUNTIME_VERSION (including latest) as satisfying the -/// minimum, which would run these tests against a sidecar that lacks the capability. -/// The tag matched here is the one the integration-tests-dapr-head CI job builds. Once -/// a dapr release ships the stateful-history protocol, this can be replaced by -/// [MinimumDaprRuntimeFact("<that version>")] and deleted. -/// -[AttributeUsage(AttributeTargets.Method, AllowMultiple = false)] -public sealed class RequiresDaprHeadFactAttribute : FactAttribute -{ - private const string RuntimeVersionEnvVarName = "DAPR_RUNTIME_VERSION"; - - /// The tag the dapr-head CI job builds and points DAPR_RUNTIME_VERSION at. - public const string DaprHeadVersion = "dapr-head"; - - /// - /// Initializes the instance. - /// - /// Populated by the compiler; forwarded so xUnit v3 can report - /// the test's source location (xUnit3003). - /// Populated by the compiler; see - /// . - public RequiresDaprHeadFactAttribute( - [CallerFilePath] string? sourceFilePath = null, - [CallerLineNumber] int sourceLineNumber = -1) - : base(sourceFilePath, sourceLineNumber) - { - var current = Environment.GetEnvironmentVariable(RuntimeVersionEnvVarName); - if (!string.Equals(current, DaprHeadVersion, StringComparison.Ordinal)) - { - Skip = $"Requires a daprd built from dapr master ({RuntimeVersionEnvVarName}=" + - $"{DaprHeadVersion}); current: '{current ?? ""}'."; - } - } -} diff --git a/test/Dapr.IntegrationTest.Workflow/StatefulHistory/StatefulHistoryTests.cs b/test/Dapr.IntegrationTest.Workflow/StatefulHistory/StatefulHistoryTests.cs index c4feb588e..16375f886 100644 --- a/test/Dapr.IntegrationTest.Workflow/StatefulHistory/StatefulHistoryTests.cs +++ b/test/Dapr.IntegrationTest.Workflow/StatefulHistory/StatefulHistoryTests.cs @@ -1,4 +1,4 @@ -// ------------------------------------------------------------------------ +// ------------------------------------------------------------------------ // Copyright 2026 The Dapr Authors // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -14,6 +14,7 @@ using Dapr.DurableTask.Protobuf; using Dapr.Testcontainers.Common; using Dapr.Testcontainers.Harnesses; +using Dapr.Testcontainers.Xunit.Attributes; using Dapr.Workflow; using Grpc.Net.ClientFactory; using Microsoft.Extensions.Configuration; @@ -26,11 +27,9 @@ namespace Dapr.IntegrationTest.Workflow.StatefulHistory; /// /// /// Requires a sidecar implementing the stateful-history protocol (dapr/durabletask-go#110, -/// reaching dapr via dapr/dapr#10142). No dapr release contains it yet, so these tests only run -/// against the image the integration-tests-dapr-head CI job builds from dapr master, and -/// skip otherwise. Against an older sidecar the capability is ignored and every turn arrives as a -/// full send, which is exactly what is written to -/// catch. +/// reaching dapr via dapr/dapr#10142), which is available starting with Dapr runtime 1.19. +/// Against an older sidecar the capability is ignored and every turn arrives as a full send, +/// which is exactly what is written to catch. /// Asserting on workflow output alone would prove nothing: a correct delta path and a sidecar /// that never sends deltas produce identical results. The counts come from a gRPC interceptor /// watching the real work-item stream. @@ -49,7 +48,7 @@ public sealed class StatefulHistoryTests private sealed record RunResult(int Deltas, int FullSends, int HistoryFetches, int Output); - [RequiresDaprHeadFact] + [MinimumDaprRuntimeFact("1.19")] public async Task DeltaDeliveryReducesFullSends() { var result = await RunAccumulateAsync(disableStatefulHistory: false); @@ -61,7 +60,7 @@ public async Task DeltaDeliveryReducesFullSends() $"expected a delta for nearly every turn: {result}"); } - [RequiresDaprHeadFact] + [MinimumDaprRuntimeFact("1.19")] public async Task WarmStreamNeverMissesItsCache() { var result = await RunAccumulateAsync(disableStatefulHistory: false); @@ -70,7 +69,7 @@ public async Task WarmStreamNeverMissesItsCache() Assert.Equal(0, result.HistoryFetches); } - [RequiresDaprHeadFact] + [MinimumDaprRuntimeFact("1.19")] public async Task DisabledWorkerReceivesOnlyFullHistories() { var result = await RunAccumulateAsync(disableStatefulHistory: true); @@ -80,6 +79,66 @@ public async Task DisabledWorkerReceivesOnlyFullHistories() Assert.True(result.FullSends >= Turns, $"expected a full send per turn: {result}"); } + [MinimumDaprRuntimeFact("1.19")] + public async Task EvictedHistoryIsRecoveredThroughGetInstanceHistory() + { + var componentsDir = TestDirectoryManager.CreateTestDirectory("stateful-history-eviction-components"); + var instanceIds = new[] { Guid.NewGuid().ToString(), Guid.NewGuid().ToString() }; + var observer = new WorkItemObserver(); + + await using var environment = await DaprTestEnvironment.CreateWithPooledNetworkAsync( + needsActorState: true, cancellationToken: TestContext.Current.CancellationToken); + await environment.StartAsync(TestContext.Current.CancellationToken); + + var harness = new DaprHarnessBuilder(componentsDir) + .WithEnvironment(environment) + .BuildWorkflow(); + + await using var testApp = await DaprHarnessBuilder.ForHarness(harness) + .ConfigureServices(builder => + { + builder.Services.AddSingleton(); + builder.Services.AddDaprWorkflowBuilder( + configureRuntime: opt => + { + opt.HistoryCacheMaxInstances = 1; + opt.RegisterWorkflow(); + opt.RegisterActivity(); + opt.RegisterActivity(); + }, + configureClient: (sp, clientBuilder) => + { + var config = sp.GetRequiredService(); + var grpcEndpoint = config["DAPR_GRPC_ENDPOINT"]; + if (!string.IsNullOrEmpty(grpcEndpoint)) + clientBuilder.UseGrpcEndpoint(grpcEndpoint); + }); + + builder.Services + .AddGrpcClient() + .AddInterceptor(() => observer); + }) + .BuildAndStartAsync(); + + using var scope = testApp.CreateScope(); + var workflowClient = scope.ServiceProvider.GetRequiredService(); + + foreach (var instanceId in instanceIds) + { + await workflowClient.ScheduleNewWorkflowAsync(nameof(EvictionWorkflow), instanceId, 0); + } + + var states = await Task.WhenAll(instanceIds.Select(instanceId => + workflowClient.WaitForWorkflowCompletionAsync( + instanceId, true, TestContext.Current.CancellationToken))); + + Assert.All(states, state => Assert.Equal(WorkflowRuntimeStatus.Completed, state.RuntimeStatus)); + Assert.All(states, state => Assert.Equal(2, state.ReadOutputAs())); + Assert.True(instanceIds.Sum(observer.Deltas) >= 2, "expected both workflows to receive a delta"); + Assert.True(instanceIds.Sum(observer.HistoryFetches) >= 1, + "expected the one-entry worker cache to evict at least one workflow history"); + } + /// /// Runs a long sequential activity chain, so each activity result is its own turn and the /// committed history grows every turn. That is what makes the omitted prefix, and therefore the @@ -161,4 +220,39 @@ public override async Task RunAsync(WorkflowContext context, int input) return current; } } + + private sealed class EvictionWorkflow : Workflow + { + public override async Task RunAsync(WorkflowContext context, int input) + { + var current = await context.CallActivityAsync(nameof(PlusOneActivity), input); + return await context.CallActivityAsync(nameof(BarrierActivity), current); + } + } + + private sealed class BarrierActivity(BarrierCoordinator coordinator) : WorkflowActivity + { + public override async Task RunAsync(WorkflowActivityContext context, int input) + { + await coordinator.SignalAndWaitAsync(); + return input + 1; + } + } + + private sealed class BarrierCoordinator + { + private readonly TaskCompletionSource _allArrived = + new(TaskCreationOptions.RunContinuationsAsynchronously); + private int _arrivals; + + public Task SignalAndWaitAsync() + { + if (Interlocked.Increment(ref _arrivals) == 2) + { + _allArrived.TrySetResult(); + } + + return _allArrived.Task.WaitAsync(TimeSpan.FromSeconds(30)); + } + } } diff --git a/test/Dapr.Workflow.Test/Worker/Grpc/GrpcProtocolHandlerStatefulHistoryTests.cs b/test/Dapr.Workflow.Test/Worker/Grpc/GrpcProtocolHandlerStatefulHistoryTests.cs index 78e5bd35f..3073aa2fe 100644 --- a/test/Dapr.Workflow.Test/Worker/Grpc/GrpcProtocolHandlerStatefulHistoryTests.cs +++ b/test/Dapr.Workflow.Test/Worker/Grpc/GrpcProtocolHandlerStatefulHistoryTests.cs @@ -1,4 +1,4 @@ -// ------------------------------------------------------------------------ +// ------------------------------------------------------------------------ // Copyright 2026 The Dapr Authors // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -41,6 +41,9 @@ private static List Events(int count) return events; } + private static List EventIds(params int[] eventIds) => + eventIds.Select(eventId => new HistoryEvent { EventId = eventId }).ToList(); + [Fact] public async Task AdvertisesStatefulHistoryCapability_ByDefault() { @@ -125,7 +128,12 @@ public async Task DeltaCacheMiss_FetchesFullHistoryViaGetInstanceHistory() var grpcClientMock = CreateGrpcClientMock(); // A delta work item whose expected prefix the (cold) cache does not hold. - var delta = new WorkflowRequest { InstanceId = "i-1", PastEvents = { Events(1) } }; + var delta = new WorkflowRequest + { + InstanceId = "i-1", + PastEvents = { EventIds(999) }, + NewEvents = { EventIds(40, 50) } + }; delta.CachedHistory = new CachedHistory { EventCount = 5 }; var workItems = new[] { new WorkItem { WorkflowRequest = delta } }; grpcClientMock @@ -136,7 +144,7 @@ public async Task DeltaCacheMiss_FetchesFullHistoryViaGetInstanceHistory() grpcClientMock .Setup(x => x.GetInstanceHistoryAsync(It.IsAny(), It.IsAny())) .Callback(() => Interlocked.Increment(ref fetchCount)) - .Returns(CreateAsyncUnaryCall(new GetInstanceHistoryResponse { Events = { Events(7) } })); + .Returns(CreateAsyncUnaryCall(new GetInstanceHistoryResponse { Events = { EventIds(10, 20, 30) } })); var completed = 0; grpcClientMock @@ -144,19 +152,22 @@ public async Task DeltaCacheMiss_FetchesFullHistoryViaGetInstanceHistory() .Callback(() => Interlocked.Increment(ref completed)) .Returns(CreateAsyncUnaryCall(new CompleteTaskResponse())); - var seenPastEvents = new ConcurrentQueue(); + IReadOnlyList? seenPastEvents = null; + IReadOnlyList? seenNewEvents = null; var handler = new GrpcProtocolHandler(grpcClientMock.Object, NullLoggerFactory.Instance); await RunHandlerUntilAsync(handler, workflowHandler: (req, _) => { - seenPastEvents.Enqueue(req.PastEvents.Count); + seenPastEvents = req.PastEvents.Select(e => e.EventId).ToArray(); + seenNewEvents = req.NewEvents.Select(e => e.EventId).ToArray(); return Task.FromResult(new WorkflowResponse { InstanceId = req.InstanceId }); }, NoActivityHandler, untilCondition: () => Volatile.Read(ref completed) >= 1, timeout: TimeSpan.FromSeconds(2)); - Assert.Equal([7], seenPastEvents); // recovered the full history via GetInstanceHistory + Assert.Equal([10, 20, 30], seenPastEvents); + Assert.Equal([40, 50], seenNewEvents); Assert.Equal(1, Volatile.Read(ref fetchCount)); } @@ -168,7 +179,7 @@ public async Task DeltaCacheHit_ReconstructsCachedPrefixPlusDelta() // Turn 1 is a full send that warms the cache; turn 2 is a delta the cache can satisfy. Turn 2 is // gated on turn 1's completion so the post-turn cache update is guaranteed to have run first. var turn1 = new WorkItem { WorkflowRequest = new WorkflowRequest { InstanceId = "i-1", PastEvents = { Events(2) } } }; - var delta = new WorkflowRequest { InstanceId = "i-1", PastEvents = { Events(1) } }; + var delta = new WorkflowRequest { InstanceId = "i-1", PastEvents = { EventIds(3) } }; delta.CachedHistory = new CachedHistory { EventCount = 2 }; var turn2 = new WorkItem { WorkflowRequest = delta }; @@ -188,19 +199,19 @@ public async Task DeltaCacheHit_ReconstructsCachedPrefixPlusDelta() }) .Returns(CreateAsyncUnaryCall(new CompleteTaskResponse())); - var seenPastEvents = new ConcurrentQueue(); + var seenPastEvents = new ConcurrentQueue(); var handler = new GrpcProtocolHandler(grpcClientMock.Object, NullLoggerFactory.Instance); await RunHandlerUntilAsync(handler, workflowHandler: (req, _) => { - seenPastEvents.Enqueue(req.PastEvents.Count); + seenPastEvents.Enqueue(req.PastEvents.Select(e => e.EventId).ToArray()); return Task.FromResult(new WorkflowResponse { InstanceId = req.InstanceId }); }, NoActivityHandler, untilCondition: () => Volatile.Read(ref completed) >= 2, timeout: TimeSpan.FromSeconds(5)); - Assert.Equal([2, 3], seenPastEvents); // turn 1 full (2); turn 2 cached prefix (2) + delta (1) + Assert.Equal([[1, 2], [1, 2, 3]], seenPastEvents); grpcClientMock.Verify( x => x.GetInstanceHistoryAsync(It.IsAny(), It.IsAny()), Times.Never); @@ -249,8 +260,13 @@ await RunHandlerUntilAsync(handler, Times.Never); } - [Fact] - public async Task CompletedWorkflow_EvictsCacheEntry_SoTheNextDeltaMisses() + [Theory] + [InlineData(OrchestrationStatus.Completed)] + [InlineData(OrchestrationStatus.Failed)] + [InlineData(OrchestrationStatus.Terminated)] + [InlineData(OrchestrationStatus.ContinuedAsNew)] + public async Task EndedWorkflow_EvictsCacheEntry_SoTheNextDeltaMisses( + OrchestrationStatus workflowStatus) { var grpcClientMock = CreateGrpcClientMock(); @@ -296,7 +312,10 @@ await RunHandlerUntilAsync(handler, // action, or the assertion below would pass for the wrong reason. if (seenPastEvents.Count == 1) { - response.Actions.Add(new WorkflowAction { CompleteWorkflow = new CompleteWorkflowAction() }); + response.Actions.Add(new WorkflowAction + { + CompleteWorkflow = new CompleteWorkflowAction { WorkflowStatus = workflowStatus } + }); } return Task.FromResult(response); @@ -308,6 +327,122 @@ await RunHandlerUntilAsync(handler, Assert.Equal(1, Volatile.Read(ref fetchCount)); } + [Fact] + public async Task RetiredStreamCannotOverwriteCurrentStreamCache() + { + var grpcClientMock = CreateGrpcClientMock(); + + var retiredTurn = new WorkItem + { + WorkflowRequest = new WorkflowRequest + { + InstanceId = "i-1", + PastEvents = { EventIds(1) } + } + }; + var failingDelta = new WorkflowRequest + { + InstanceId = "force-reconnect", + PastEvents = { EventIds(2) }, + CachedHistory = new CachedHistory { EventCount = 5 } + }; + var currentTurn = new WorkItem + { + WorkflowRequest = new WorkflowRequest + { + InstanceId = "i-1", + PastEvents = { EventIds(10, 11) } + } + }; + var currentDelta = new WorkflowRequest + { + InstanceId = "i-1", + PastEvents = { EventIds(12) }, + CachedHistory = new CachedHistory { EventCount = 2 } + }; + + var neverCompletes = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var retiredHandlerStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var releaseRetiredHandler = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var retiredCompletionSent = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + grpcClientMock + .SetupSequence(x => x.GetWorkItems(It.IsAny(), It.IsAny())) + .Returns(CreateServerStreamingCallFromReader(new GatedStreamReader( + [ + (retiredTurn, null), + (new WorkItem { WorkflowRequest = failingDelta }, null), + (new WorkItem(), neverCompletes.Task) + ]))) + .Returns(CreateServerStreamingCallFromReader(new GatedStreamReader( + [ + (currentTurn, null), + (new WorkItem { WorkflowRequest = currentDelta }, retiredCompletionSent.Task) + ]))) + .Returns(CreateServerStreamingCall([])); + + var unexpectedCurrentStreamFetches = 0; + grpcClientMock + .Setup(x => x.GetInstanceHistoryAsync(It.IsAny(), It.IsAny())) + .Returns((request, _) => + { + if (request.InstanceId == "force-reconnect") + { + throw new RpcException(new Status(StatusCode.Unavailable, "force reconnect")); + } + + Interlocked.Increment(ref unexpectedCurrentStreamFetches); + return CreateAsyncUnaryCall(new GetInstanceHistoryResponse { Events = { EventIds(90, 91) } }); + }); + + var completed = 0; + grpcClientMock + .Setup(x => x.CompleteOrchestratorTaskAsync(It.IsAny(), It.IsAny())) + .Callback((response, _) => + { + Interlocked.Increment(ref completed); + if (response.CustomStatus == "current") + { + releaseRetiredHandler.TrySetResult(); + } + else if (response.CustomStatus == "retired") + { + retiredCompletionSent.TrySetResult(); + } + }) + .Returns(CreateAsyncUnaryCall(new CompleteTaskResponse())); + + IReadOnlyList? reconstructedHistory = null; + var handler = new GrpcProtocolHandler(grpcClientMock.Object, NullLoggerFactory.Instance); + + await RunHandlerUntilAsync(handler, + workflowHandler: async (request, _) => + { + var eventIds = request.PastEvents.Select(e => e.EventId).ToArray(); + if (eventIds.SequenceEqual([1])) + { + retiredHandlerStarted.TrySetResult(); + await releaseRetiredHandler.Task; + return new WorkflowResponse { InstanceId = request.InstanceId, CustomStatus = "retired" }; + } + + if (eventIds.SequenceEqual([10, 11])) + { + await retiredHandlerStarted.Task; + return new WorkflowResponse { InstanceId = request.InstanceId, CustomStatus = "current" }; + } + + reconstructedHistory = eventIds; + return new WorkflowResponse { InstanceId = request.InstanceId, CustomStatus = "delta" }; + }, + NoActivityHandler, + untilCondition: () => Volatile.Read(ref completed) >= 3, + timeout: TimeSpan.FromSeconds(20)); + + Assert.Equal([10, 11, 12], reconstructedHistory); + Assert.Equal(0, Volatile.Read(ref unexpectedCurrentStreamFetches)); + } + [Fact] public async Task Reconnect_ResetsCache_SoTheNextDeltaMisses() { diff --git a/test/Dapr.Workflow.Test/Worker/Grpc/WorkflowHistoryCacheTests.cs b/test/Dapr.Workflow.Test/Worker/Grpc/WorkflowHistoryCacheTests.cs index ce25f3689..632e26955 100644 --- a/test/Dapr.Workflow.Test/Worker/Grpc/WorkflowHistoryCacheTests.cs +++ b/test/Dapr.Workflow.Test/Worker/Grpc/WorkflowHistoryCacheTests.cs @@ -1,4 +1,4 @@ -// ------------------------------------------------------------------------ +// ------------------------------------------------------------------------ // Copyright 2026 The Dapr Authors // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -110,6 +110,21 @@ public void CountCapEvictsLeastRecentlyUsed() Assert.NotNull(cache.Get("c")); } + [Fact] + public void GetMovesEntryToMostRecentlyUsedPosition() + { + var cache = new WorkflowHistoryCache(maxInstances: 2); + + cache.Put("a", Events(1), cache.Generation); + cache.Put("b", Events(1), cache.Generation); + Assert.NotNull(cache.Get("a")); + cache.Put("c", Events(1), cache.Generation); + + Assert.NotNull(cache.Get("a")); + Assert.Null(cache.Get("b")); + Assert.NotNull(cache.Get("c")); + } + [Fact] public void ByteCapEvictsLeastRecentlyUsed() { @@ -189,12 +204,11 @@ public void TtlSweepIsSliding() } [Fact] - public void NonPositiveConfigUsesDefaults() + public void ZeroConfigUsesDefaults() { - // ttl/maxInstances fall back to their (large) defaults; maxBytes becomes unlimited. None of these + // maxInstances falls back to its large default and maxBytes becomes unlimited. Neither // should evict the three modest entries below. - var cache = new WorkflowHistoryCache( - ttl: TimeSpan.Zero, maxInstances: -1, maxBytes: -5); + var cache = new WorkflowHistoryCache(maxInstances: 0, maxBytes: 0); cache.Put("a", Events(1), cache.Generation); cache.Put("b", Events(1), cache.Generation); @@ -205,4 +219,68 @@ public void NonPositiveConfigUsesDefaults() Assert.NotNull(cache.Get("b")); Assert.NotNull(cache.Get("c")); } + + [Fact] + public void InvalidConfigThrows() + { + Assert.Throws(() => new WorkflowHistoryCache(ttl: TimeSpan.Zero)); + Assert.Throws(() => new WorkflowHistoryCache(ttl: TimeSpan.FromSeconds(-1))); + Assert.Throws(() => new WorkflowHistoryCache(maxInstances: -1)); + Assert.Throws(() => new WorkflowHistoryCache(maxBytes: -1)); + } + + [Fact] + public async Task ConcurrentOperationsPreserveBoundsAndAccounting() + { + const int maxInstances = 16; + var entryBytes = BytesOf(3); + var cache = new WorkflowHistoryCache( + maxInstances: maxInstances, + maxBytes: entryBytes * maxInstances); + + var tasks = Enumerable.Range(0, 8).Select(worker => Task.Run(() => + { + for (var i = 0; i < 500; i++) + { + var instanceId = $"instance-{(worker * 500 + i) % 64}"; + var generation = cache.Generation; + cache.Put(instanceId, Events(3), generation); + _ = cache.Get(instanceId); + + if (i % 7 == 0) + { + cache.Remove($"instance-{(i + 1) % 64}", generation); + } + } + })); + + await Task.WhenAll(tasks); + + Assert.InRange(cache.Count, 0, maxInstances); + Assert.InRange(cache.TotalBytes, 0, entryBytes * maxInstances); + } + + [Fact] + public async Task ConcurrentResetRejectsRetiredWrites() + { + var cache = new WorkflowHistoryCache(); + var retiredGeneration = cache.Generation; + using var start = new ManualResetEventSlim(); + + var writers = Enumerable.Range(0, 8).Select(worker => Task.Run(() => + { + start.Wait(); + for (var i = 0; i < 250; i++) + { + cache.Put($"retired-{worker}-{i}", Events(1), retiredGeneration); + } + })).ToArray(); + + cache.Reset(); + start.Set(); + await Task.WhenAll(writers); + + Assert.Equal(0, cache.Count); + Assert.Equal(0, cache.TotalBytes); + } } diff --git a/test/Dapr.Workflow.Test/Worker/WorkflowWorkerTests.cs b/test/Dapr.Workflow.Test/Worker/WorkflowWorkerTests.cs index 8a6730097..c2b7b18e3 100644 --- a/test/Dapr.Workflow.Test/Worker/WorkflowWorkerTests.cs +++ b/test/Dapr.Workflow.Test/Worker/WorkflowWorkerTests.cs @@ -1,4 +1,4 @@ -// ------------------------------------------------------------------------ +// ------------------------------------------------------------------------ // Copyright 2025 The Dapr Authors // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -115,7 +115,7 @@ public async Task StopAsync_ShouldDisposeProtocolHandler_WhenPresent() await worker.StopAsync(CancellationToken.None); } - + [Fact] public async Task ExecuteAsync_ShouldComplete_WhenGrpcStreamCompletesImmediately() { @@ -249,7 +249,7 @@ public async Task HandleWorkflowResponseAsync_ShouldReturnTerminatedCompletion_W Assert.NotNull(action.CompleteWorkflow); Assert.Equal(OrchestrationStatus.Terminated, action.CompleteWorkflow!.WorkflowStatus); } - + [Fact] public async Task HandleWorkflowResponseAsync_ShouldNotReturnTerminatedCompletion_WhenReplayLatestEventIsNotExecutionTerminated() { @@ -296,7 +296,7 @@ public async Task HandleWorkflowResponseAsync_ShouldNotReturnTerminatedCompletio Assert.NotEqual(OrchestrationStatus.Terminated, action.CompleteWorkflow!.WorkflowStatus); Assert.Equal(OrchestrationStatus.Failed, action.CompleteWorkflow.WorkflowStatus); } - + [Fact] public async Task HandleWorkflowResponseAsync_ShouldReturnEmptyResponse_WhenLatestEventIsExecutionSuspended() { @@ -334,7 +334,7 @@ public async Task HandleWorkflowResponseAsync_ShouldReturnEmptyResponse_WhenLate Assert.Equal("i", response.InstanceId); Assert.Empty(response.Actions); } - + [Fact] public async Task HandleWorkflowResponseAsync_ShouldNotShortCircuit_WhenLatestEventIsExecutionResumed() { @@ -462,7 +462,7 @@ public void CreateCallOptions_ShouldNotIncludeApiTokenHeader_WhenTokenIsEmpty() Assert.False(HasHeader(callOptions, "dapr-api-token", out _)); Assert.True(HasHeader(callOptions, "User-Agent", out _)); } - + [Fact] public async Task CallChildWorkflowAsync_ShouldComplete_WhenCompletionEventArrivesLater() { @@ -510,13 +510,13 @@ public async Task CallChildWorkflowAsync_ShouldComplete_WhenCompletionEventArriv Assert.Equal(99, value); Assert.Empty(context.PendingActions); } - + [Fact] public void CallChildWorkflowAsync_ShouldPreserveRouterTargetAppId_OnScheduledAction() { const string appId1 = "this-app"; const string appId2 = "remote-app"; - + var serializer = new JsonDaprSerializer(new JsonSerializerOptions(JsonSerializerDefaults.Web)); var context = new WorkflowOrchestrationContext( "wf", "parent", new DateTime(2025, 01, 01, 0, 0, 0, DateTimeKind.Utc), @@ -530,7 +530,7 @@ public void CallChildWorkflowAsync_ShouldPreserveRouterTargetAppId_OnScheduledAc Assert.Equal(appId1, action.Router.SourceAppID); Assert.Equal(appId2, action.Router.TargetAppID); } - + [Fact] public async Task CallChildWorkflowAsync_ShouldComplete_WhenCompletionArrivedBeforeCall() { @@ -570,7 +570,7 @@ public async Task CallChildWorkflowAsync_ShouldComplete_WhenCompletionArrivedBef Assert.Equal(13, value); Assert.Empty(context.PendingActions); } - + [Fact] public async Task CallChildWorkflowAsync_ShouldIgnoreDuplicateCompletionEvents() { @@ -610,7 +610,7 @@ public async Task CallChildWorkflowAsync_ShouldIgnoreDuplicateCompletionEvents() var value = await task; Assert.Equal(200, value); } - + [Fact] public async Task HandleWorkflowResponseAsync_ShouldAllowWorkflowToComplete_OnSecondPass_WhenChildCompletionInHistory() { @@ -687,7 +687,7 @@ await ctx.CallChildWorkflowAsync("TargetWorkflow", input: 7, Assert.Contains(resp2.Actions, a => a.CompleteWorkflow != null); Assert.Equal(OrchestrationStatus.Completed, resp2.Actions.Single(a => a.CompleteWorkflow != null).CompleteWorkflow!.WorkflowStatus); } - + [Fact] public async Task CallChildWorkflowAsync_ShouldOnlyCompleteAfterCreation_WhenCompletionArrivesFirst() { @@ -721,7 +721,7 @@ public async Task CallChildWorkflowAsync_ShouldOnlyCompleteAfterCreation_WhenCom var value = await task; Assert.Equal(21, value); } - + [Fact] public async Task CallChildWorkflowAsync_ShouldCompleteOnlyForMatchingTaskScheduledId_WhenReplaySchedulesAgain() { @@ -914,7 +914,7 @@ public async Task HandleWorkflowResponseAsync_ShouldCompleteWorkflow_AndIncludeO var completion = response.Actions .FirstOrDefault(a => a.CompleteWorkflow != null)?.CompleteWorkflow; - + Assert.NotNull(completion); Assert.Equal(OrchestrationStatus.Completed, completion.WorkflowStatus); Assert.Equal("42", completion.Result); @@ -1293,7 +1293,7 @@ public async Task HandleOrchestratorResponseAsync_ShouldKeepTraceContext_ForWork Assert.Equal(expectedTraceId, logProvider.GetTraceIdForMessage("workflow-user-log")); Assert.Equal(expectedTraceId, logProvider.GetTraceIdForMessage("Workflow execution completed")); } - + [Fact] public async Task ExecuteAsync_ShouldRetry_WhenGrpcProtocolHandlerStartFailsWithException() { @@ -1401,7 +1401,7 @@ public async Task HandleWorkflowResponseAsync_ShouldAdvanceCurrentUtcDateTime_Wh Assert.Equal(beginDateTime, ctx.CurrentUtcDateTime); await ctx.CreateTimer(TimeSpan.FromSeconds(5)); Assert.Equal(beginDateTime.AddSeconds(5), ctx.CurrentUtcDateTime); - + return null; })); @@ -1626,7 +1626,7 @@ public async Task HandleWorkflowResponseAsync_ShouldCompleted_WhenEventReceived( } } }; - + var response = await InvokeHandleWorkflowResponseAsync(worker, request); @@ -1704,7 +1704,7 @@ public async Task HandleWorkflowResponseAsync_ShouldReturnFailureDetails_WhenTim new HistoryEvent { EventRaised = new EventRaisedEvent { Name = "myevent" } } } }; - + var response = await InvokeHandleWorkflowResponseAsync(worker, request); @@ -1747,62 +1747,62 @@ public async Task HandleActivityResponseAsync_ShouldUseEmptyInstanceId_WhenWorkf Assert.Null(response.FailureDetails); Assert.Equal(string.Empty, response.Result); } - + // ------------------------------------------------------------------------- // RequiresHistoryStreaming // ------------------------------------------------------------------------- - // [Fact] - // public async Task HandleWorkflowResponseAsync_ShouldStreamHistory_WhenRequiresHistoryStreamingIsTrue() - // { - // // When RequiresHistoryStreaming is set, the worker must fetch past history - // // via StreamInstanceHistory and merge it with the inline PastEvents before - // // running the workflow. Here we put the ExecutionStarted event inside the - // // stream (not in PastEvents) so the workflow can only complete if streaming works. - // var sp = new ServiceCollection().BuildServiceProvider(); - // var serializer = new JsonDaprSerializer(new JsonSerializerOptions(JsonSerializerDefaults.Web)); - // - // var factory = new StubWorkflowsFactory(); - // factory.AddWorkflow("wf", new InlineWorkflow( - // inputType: typeof(int), - // run: (_, input) => Task.FromResult((int)input! + 1))); - // - // // The streamed chunk carries the ExecutionStarted event. - // var streamedChunk = new HistoryChunk(); - // streamedChunk.Events.Add(new HistoryEvent - // { - // ExecutionStarted = new ExecutionStartedEvent { Name = "wf", Input = "10" } - // }); - // - // var grpcClientMock = CreateGrpcClientMock(); - // grpcClientMock - // .Setup(x => x.GetInstanceHistoryAsync(It.IsAny(), It.IsAny())) - // .Returns(CreateHistoryStreamingCall(SingleItemAsync(streamedChunk))); - // - // var worker = new WorkflowWorker( - // grpcClientMock.Object, - // factory, - // NullLoggerFactory.Instance, - // serializer, - // sp); - // - // var request = new WorkflowRequest - // { - // InstanceId = "stream-i", - // RequiresHistoryStreaming = true - // }; - // - // var response = await InvokeHandleWorkflowResponseAsync(worker, request); - // - // Assert.Equal("stream-i", response.InstanceId); - // var complete = response.Actions.Single(a => a.CompleteWorkflow != null).CompleteWorkflow!; - // Assert.Equal(OrchestrationStatus.Completed, complete.WorkflowStatus); - // Assert.Equal("11", complete.Result); - // - // grpcClientMock.Verify( - // x => x.StreamInstanceHistory(It.IsAny(), It.IsAny()), - // Times.Once()); - // } + [Fact] + public async Task HandleWorkflowResponseAsync_ShouldFetchHistory_WhenRequiresHistoryStreamingIsTrue() + { + var sp = new ServiceCollection().BuildServiceProvider(); + var serializer = new JsonDaprSerializer(new JsonSerializerOptions(JsonSerializerDefaults.Web)); + + var factory = new StubWorkflowsFactory(); + factory.AddWorkflow("wf", new InlineWorkflow( + inputType: typeof(int), + run: (_, input) => Task.FromResult((int)input! + 1))); + + var grpcClientMock = CreateGrpcClientMock(); + grpcClientMock + .Setup(x => x.GetInstanceHistoryAsync( + It.Is(request => request.InstanceId == "stream-i"), + It.IsAny())) + .Returns(CreateAsyncUnaryCall(new GetInstanceHistoryResponse + { + Events = + { + new HistoryEvent + { + ExecutionStarted = new ExecutionStartedEvent { Name = "wf", Input = "10" } + } + } + })); + + var worker = new WorkflowWorker( + grpcClientMock.Object, + factory, + NullLoggerFactory.Instance, + serializer, + sp); + + var response = await InvokeHandleWorkflowResponseAsync(worker, new WorkflowRequest + { + InstanceId = "stream-i", + RequiresHistoryStreaming = true + }); + + Assert.Equal("stream-i", response.InstanceId); + var complete = response.Actions.Single(a => a.CompleteWorkflow != null).CompleteWorkflow!; + Assert.Equal(OrchestrationStatus.Completed, complete.WorkflowStatus); + Assert.Equal("11", complete.Result); + + grpcClientMock.Verify( + x => x.GetInstanceHistoryAsync( + It.Is(request => request.InstanceId == "stream-i"), + It.IsAny()), + Times.Once()); + } // ------------------------------------------------------------------------- // Workflow-name extraction fallbacks @@ -2480,6 +2480,9 @@ private static async Task InvokeHandleActivityResponseAsync(Wo return await task; } + private static AsyncUnaryCall CreateAsyncUnaryCall(T response) => + new(Task.FromResult(response), Task.FromResult(new Metadata()), () => Status.DefaultSuccess, () => [], () => { }); + private static Mock CreateGrpcClientMock() { var callInvoker = new Mock(MockBehavior.Loose); diff --git a/test/Dapr.Workflow.Test/WorkflowRuntimeOptionsTests.cs b/test/Dapr.Workflow.Test/WorkflowRuntimeOptionsTests.cs index 5d06df5fc..8f3f523ea 100644 --- a/test/Dapr.Workflow.Test/WorkflowRuntimeOptionsTests.cs +++ b/test/Dapr.Workflow.Test/WorkflowRuntimeOptionsTests.cs @@ -21,6 +21,54 @@ namespace Dapr.Workflow.Test; public class WorkflowRuntimeOptionsTests { + [Theory] + [InlineData(-1)] + [InlineData(-1024)] + public void HistoryCacheMaxInstances_ShouldRejectNegativeValues(int value) + { + var options = new WorkflowRuntimeOptions(); + + Assert.Throws(() => options.HistoryCacheMaxInstances = value); + } + + [Theory] + [InlineData(-1)] + [InlineData(-1024)] + public void HistoryCacheMaxBytes_ShouldRejectNegativeValues(long value) + { + var options = new WorkflowRuntimeOptions(); + + Assert.Throws(() => options.HistoryCacheMaxBytes = value); + } + + [Fact] + public void HistoryCacheTtl_ShouldRejectNonPositiveValues() + { + var options = new WorkflowRuntimeOptions(); + + Assert.Throws(() => options.HistoryCacheTtl = TimeSpan.Zero); + Assert.Throws(() => options.HistoryCacheTtl = TimeSpan.FromSeconds(-1)); + } + + [Fact] + public void HistoryCacheOptions_ShouldAcceptDocumentedDefaultsAndPositiveValues() + { + var options = new WorkflowRuntimeOptions + { + HistoryCacheTtl = TimeSpan.FromMinutes(5), + HistoryCacheMaxInstances = 10, + HistoryCacheMaxBytes = 1024 + }; + + Assert.Equal(TimeSpan.FromMinutes(5), options.HistoryCacheTtl); + Assert.Equal(10, options.HistoryCacheMaxInstances); + Assert.Equal(1024, options.HistoryCacheMaxBytes); + + options.HistoryCacheTtl = null; + options.HistoryCacheMaxInstances = 0; + options.HistoryCacheMaxBytes = 0; + } + [Fact] public void UseGrpcChannelOptions_ShouldThrowArgumentNullException_WhenNull() {