diff --git a/src/Dapr.Actors/Communication/ActorStateResponse.cs b/src/Dapr.Actors/Communication/ActorStateResponse.cs index bf48467e7..7af34603f 100644 --- a/src/Dapr.Actors/Communication/ActorStateResponse.cs +++ b/src/Dapr.Actors/Communication/ActorStateResponse.cs @@ -25,10 +25,12 @@ public class ActorStateResponse /// /// The response value. /// The time to live expiration time. - public ActorStateResponse(T value, DateTimeOffset? ttlExpireTime) + /// The HTTP status code provided alongside the state response. + public ActorStateResponse(T value, DateTimeOffset? ttlExpireTime, int httpStatusCode = 200) { this.Value = value; this.TTLExpireTime = ttlExpireTime; + this.HttpStatusCode = httpStatusCode; } /// @@ -46,4 +48,9 @@ public ActorStateResponse(T value, DateTimeOffset? ttlExpireTime) /// The time to live expiration time. /// public DateTimeOffset? TTLExpireTime { get; } -} \ No newline at end of file + + /// + /// The HTTP status code provided alongside the state response. + /// + public int HttpStatusCode { get; } +} diff --git a/src/Dapr.Actors/DaprHttpInteractor.cs b/src/Dapr.Actors/DaprHttpInteractor.cs index b7b1f53a6..2c30a4979 100644 --- a/src/Dapr.Actors/DaprHttpInteractor.cs +++ b/src/Dapr.Actors/DaprHttpInteractor.cs @@ -32,7 +32,7 @@ namespace Dapr.Actors; /// /// Class to interact with Dapr runtime over http. /// -internal class DaprHttpInteractor : IDaprInteractor +internal sealed class DaprHttpInteractor : IDaprInteractor { private readonly JsonSerializerOptions jsonSerializerOptions = JsonSerializerDefaults.Web; private readonly string httpEndpoint; @@ -71,7 +71,7 @@ HttpRequestMessage RequestFunc() } using var response = await this.SendAsync(RequestFunc, relativeUrl, cancellationToken); - var stringResponse = await response.Content.ReadAsStringAsync(); + var stringResponse = await response.Content.ReadAsStringAsync(cancellationToken); DateTimeOffset? ttlExpireTime = null; if (response.Headers.TryGetValues(Constants.TTLResponseHeaderName, out IEnumerable headerValues)) @@ -83,13 +83,15 @@ HttpRequestMessage RequestFunc() } } - return new ActorStateResponse(stringResponse, ttlExpireTime); + return new ActorStateResponse(stringResponse, ttlExpireTime, (int)response.StatusCode); } public Task SaveStateTransactionallyAsync(string actorType, string actorId, string data, CancellationToken cancellationToken = default) { var relativeUrl = string.Format(CultureInfo.InvariantCulture, Constants.ActorStateRelativeUrlFormat, actorType, actorId); + return this.SendAsync(RequestFunc, relativeUrl, cancellationToken); + HttpRequestMessage RequestFunc() { var request = new HttpRequestMessage() @@ -100,8 +102,6 @@ HttpRequestMessage RequestFunc() return request; } - - return this.SendAsync(RequestFunc, relativeUrl, cancellationToken); } public async Task InvokeActorMethodWithRemotingAsync(ActorMessageSerializersManager serializersManager, IActorRequestMessage remotingRequestRequestMessage, CancellationToken cancellationToken = default) @@ -116,8 +116,8 @@ public async Task InvokeActorMethodWithRemotingAsync(Acto var serializedHeader = serializersManager.GetHeaderSerializer() .SerializeRequestHeader(remotingRequestRequestMessage.GetHeader()); - var msgBodySeriaizer = serializersManager.GetRequestMessageBodySerializer(interfaceId, methodName); - var serializedMsgBody = msgBodySeriaizer.Serialize(remotingRequestRequestMessage.GetBody()); + var msgBodySerializer = serializersManager.GetRequestMessageBodySerializer(interfaceId, methodName); + var serializedMsgBody = msgBodySerializer.Serialize(remotingRequestRequestMessage.GetBody()); // Send Request var relativeUrl = string.Format(CultureInfo.InvariantCulture, Constants.ActorMethodRelativeUrlFormat, actorType, actorId, methodName); @@ -166,7 +166,7 @@ HttpRequestMessage RequestFunc() IActorResponseMessageBody actorResponseMessageBody = null; if (retval != null && retval.Content != null) { - var responseMessageBody = await retval.Content.ReadAsStreamAsync(); + var responseMessageBody = await retval.Content.ReadAsStreamAsync(cancellationToken); // Deserialize Actor Response Message Body // Deserialize to ActorInvokeException when there is response header otherwise normal path @@ -242,7 +242,7 @@ HttpRequestMessage RequestFunc() } var response = await this.SendAsync(RequestFunc, relativeUrl, cancellationToken); - var stream = await response.Content.ReadAsStreamAsync(); + var stream = await response.Content.ReadAsStreamAsync(cancellationToken); return stream; } @@ -364,9 +364,9 @@ internal async Task SendAsyncGetResponseAsRawJson( using var response = await this.SendAsyncHandleUnsuccessfulResponse(requestFunc, relativeUri, cancellationToken); var retValue = default(string); - if (response != null && response.Content != null) + if (response?.Content != null) { - retValue = await response.Content.ReadAsStringAsync(); + retValue = await response.Content.ReadAsStringAsync(cancellationToken); } return retValue; @@ -376,7 +376,7 @@ internal async Task SendAsyncGetResponseAsRawJson( /// Disposes resources. /// /// False values indicates the method is being called by the runtime, true value indicates the method is called by the user code. - protected virtual void Dispose(bool disposing) + private void Dispose(bool disposing) { if (!this.disposed) { @@ -425,10 +425,10 @@ HttpRequestMessage FinalRequestFunc() try { - var contentStream = await response.Content.ReadAsStreamAsync(); + var contentStream = await response.Content.ReadAsStreamAsync(cancellationToken); if (contentStream.Length != 0) { - error = await JsonSerializer.DeserializeAsync(contentStream, jsonSerializerOptions); + error = await JsonSerializer.DeserializeAsync(contentStream, jsonSerializerOptions, cancellationToken); } } catch (Exception ex) @@ -505,7 +505,6 @@ private void AddDaprApiTokenHeader(HttpRequestMessage request) if (!string.IsNullOrWhiteSpace(this.daprApiToken)) { request.Headers.Add("dapr-api-token", this.daprApiToken); - return; } } -} \ No newline at end of file +} diff --git a/src/Dapr.Actors/Runtime/ActorStateManager.cs b/src/Dapr.Actors/Runtime/ActorStateManager.cs index 6ebfd3271..8735fc917 100644 --- a/src/Dapr.Actors/Runtime/ActorStateManager.cs +++ b/src/Dapr.Actors/Runtime/ActorStateManager.cs @@ -86,9 +86,14 @@ public async Task TryAddStateAsync(string stateName, T value, Cancellat var stateChangeTracker = GetContextualStateTracker(); - if (stateChangeTracker.ContainsKey(stateName)) + if (stateChangeTracker.TryGetValue(stateName, out var stateMetadata)) { - var stateMetadata = stateChangeTracker[stateName]; + // The state is known to not exist in the state store, so it can be added without consulting it again + if (stateMetadata.ChangeKind == StateChangeKind.NotFound) + { + stateChangeTracker[stateName] = StateMetadata.Create(value, StateChangeKind.Add); + return true; + } // Check if the property was marked as remove or is expired in the cache if (stateMetadata.ChangeKind == StateChangeKind.Remove || (stateMetadata.TTLExpireTime.HasValue && stateMetadata.TTLExpireTime.Value <= DateTimeOffset.UtcNow)) @@ -117,9 +122,14 @@ public async Task TryAddStateAsync(string stateName, T value, TimeSpan var stateChangeTracker = GetContextualStateTracker(); - if (stateChangeTracker.ContainsKey(stateName)) + if (stateChangeTracker.TryGetValue(stateName, out var stateMetadata)) { - var stateMetadata = stateChangeTracker[stateName]; + // The state is known to not exist in the state store, so it can be added without consulting it again + if (stateMetadata.ChangeKind == StateChangeKind.NotFound) + { + stateChangeTracker[stateName] = StateMetadata.Create(value, StateChangeKind.Add, ttl: ttl); + return true; + } // Check if the property was marked as remove in the cache or has been expired. if (stateMetadata.ChangeKind == StateChangeKind.Remove || (stateMetadata.TTLExpireTime.HasValue && stateMetadata.TTLExpireTime.Value <= DateTimeOffset.UtcNow)) @@ -162,12 +172,10 @@ public async Task> TryGetStateAsync(string stateName, Can var stateChangeTracker = GetContextualStateTracker(); - if (stateChangeTracker.ContainsKey(stateName)) + if (stateChangeTracker.TryGetValue(stateName, out var stateMetadata)) { - var stateMetadata = stateChangeTracker[stateName]; - - // Check if the property was marked as remove in the cache or is expired - if (stateMetadata.ChangeKind == StateChangeKind.Remove || (stateMetadata.TTLExpireTime.HasValue && stateMetadata.TTLExpireTime.Value <= DateTimeOffset.UtcNow)) + // Check if the property wasn't populated, was marked as remove in the cache or is expired + if (stateMetadata.ChangeKind == StateChangeKind.NotFound || stateMetadata.ChangeKind == StateChangeKind.Remove || (stateMetadata.TTLExpireTime.HasValue && stateMetadata.TTLExpireTime.Value <= DateTimeOffset.UtcNow)) { return new ConditionalValue(false, default); } @@ -175,13 +183,17 @@ public async Task> TryGetStateAsync(string stateName, Can return new ConditionalValue(true, (T)stateMetadata.Value); } - var conditionalResult = await this.TryGetStateFromStateProviderAsync(stateName, cancellationToken); - if (conditionalResult.HasValue) + var (response, httpStatusCode) = await this.TryGetStateFromStateProviderAsync(stateName, cancellationToken); + if (httpStatusCode == 200 && response is not null) // HTTP status code == 200 OK { - stateChangeTracker.Add(stateName, StateMetadata.Create(conditionalResult.Value.Value, StateChangeKind.None, ttlExpireTime: conditionalResult.Value.TTLExpireTime)); - return new ConditionalValue(true, conditionalResult.Value.Value); + stateChangeTracker.Add(stateName, StateMetadata.Create(response.Value, StateChangeKind.None, ttlExpireTime: response.TTLExpireTime)); + return new ConditionalValue(true, response.Value); } + // Value isn't forthcoming at all + // Only other expected status codes are 204 (key not found, empty response), 400 (actor not found) and 500 (request failed) all of which are not + // transient value responses, so cache the absence of the value so subsequent lookups don't go back to the state store. + stateChangeTracker[stateName] = StateMetadata.CreateNotFound(); return new ConditionalValue(false, default); } @@ -193,14 +205,19 @@ public async Task SetStateAsync(string stateName, T value, CancellationToken var stateChangeTracker = GetContextualStateTracker(); - if (stateChangeTracker.ContainsKey(stateName)) + if (stateChangeTracker.TryGetValue(stateName, out var stateMetadata)) { - var stateMetadata = stateChangeTracker[stateName]; + // The state is known to not exist in the state store, so this is an add rather than an update + if (stateMetadata.ChangeKind == StateChangeKind.NotFound) + { + stateChangeTracker[stateName] = StateMetadata.Create(value, StateChangeKind.Add); + return; + } + stateMetadata.Value = value; stateMetadata.TTLExpireTime = null; - if (stateMetadata.ChangeKind == StateChangeKind.None || - stateMetadata.ChangeKind == StateChangeKind.Remove) + if (stateMetadata.ChangeKind is StateChangeKind.None or StateChangeKind.Remove) { stateMetadata.ChangeKind = StateChangeKind.Update; } @@ -223,14 +240,19 @@ public async Task SetStateAsync(string stateName, T value, TimeSpan ttl, Canc var stateChangeTracker = GetContextualStateTracker(); - if (stateChangeTracker.ContainsKey(stateName)) + if (stateChangeTracker.TryGetValue(stateName, out var stateMetadata)) { - var stateMetadata = stateChangeTracker[stateName]; + // The state is known to not exist in the state store, so this is an add rather than an update + if (stateMetadata.ChangeKind == StateChangeKind.NotFound) + { + stateChangeTracker[stateName] = StateMetadata.Create(value, StateChangeKind.Add, ttl: ttl); + return; + } + stateMetadata.Value = value; stateMetadata.TTLExpireTime = DateTimeOffset.UtcNow.Add(ttl); - if (stateMetadata.ChangeKind == StateChangeKind.None || - stateMetadata.ChangeKind == StateChangeKind.Remove) + if (stateMetadata.ChangeKind is StateChangeKind.None or StateChangeKind.Remove) { stateMetadata.ChangeKind = StateChangeKind.Update; } @@ -263,9 +285,13 @@ public async Task TryRemoveStateAsync(string stateName, CancellationToken var stateChangeTracker = GetContextualStateTracker(); - if (stateChangeTracker.ContainsKey(stateName)) + if (stateChangeTracker.TryGetValue(stateName, out var stateMetadata)) { - var stateMetadata = stateChangeTracker[stateName]; + // The state is known to not exist in the state store, so there's nothing to remove + if (stateMetadata.ChangeKind == StateChangeKind.NotFound) + { + return false; + } if (stateMetadata.TTLExpireTime.HasValue && stateMetadata.TTLExpireTime.Value <= DateTimeOffset.UtcNow) { @@ -303,12 +329,10 @@ public async Task ContainsStateAsync(string stateName, CancellationToken c var stateChangeTracker = GetContextualStateTracker(); - if (stateChangeTracker.ContainsKey(stateName)) + if (stateChangeTracker.TryGetValue(stateName, out var stateMetadata)) { - var stateMetadata = stateChangeTracker[stateName]; - - // Check if the property was marked as remove in the cache - return stateMetadata.ChangeKind != StateChangeKind.Remove; + // Check if the property is known to not exist or was marked as remove in the cache + return stateMetadata.ChangeKind is not (StateChangeKind.Remove or StateChangeKind.NotFound); } if (await this.actor.Host.StateProvider.ContainsStateAsync(this.actorTypeName, this.actor.Id.ToString(), stateName, cancellationToken)) @@ -367,9 +391,14 @@ public async Task AddOrUpdateStateAsync( var stateChangeTracker = GetContextualStateTracker(); - if (stateChangeTracker.ContainsKey(stateName)) + if (stateChangeTracker.TryGetValue(stateName, out var stateMetadata)) { - var stateMetadata = stateChangeTracker[stateName]; + // The state is known to not exist in the state store, so add it without consulting the store again + if (stateMetadata.ChangeKind == StateChangeKind.NotFound) + { + stateChangeTracker[stateName] = StateMetadata.Create(addValue, StateChangeKind.Add); + return addValue; + } // Check if the property was marked as remove in the cache if (stateMetadata.ChangeKind == StateChangeKind.Remove) @@ -389,10 +418,10 @@ public async Task AddOrUpdateStateAsync( return newValue; } - var conditionalResult = await this.TryGetStateFromStateProviderAsync(stateName, cancellationToken); - if (conditionalResult.HasValue) + var (response, httpStatusCode) = await this.TryGetStateFromStateProviderAsync(stateName, cancellationToken); + if (httpStatusCode == 200 && response is not null) { - var newValue = updateValueFactory.Invoke(stateName, conditionalResult.Value.Value); + var newValue = updateValueFactory.Invoke(stateName, response.Value); stateChangeTracker.Add(stateName, StateMetadata.Create(newValue, StateChangeKind.Update)); return newValue; @@ -415,9 +444,14 @@ public async Task AddOrUpdateStateAsync( var stateChangeTracker = GetContextualStateTracker(); - if (stateChangeTracker.ContainsKey(stateName)) + if (stateChangeTracker.TryGetValue(stateName, out var stateMetadata)) { - var stateMetadata = stateChangeTracker[stateName]; + // The state is known to not exist in the state store, so add it without consulting the store again + if (stateMetadata.ChangeKind == StateChangeKind.NotFound) + { + stateChangeTracker[stateName] = StateMetadata.Create(addValue, StateChangeKind.Add, ttl: ttl); + return addValue; + } // Check if the property was marked as remove in the cache if (stateMetadata.ChangeKind == StateChangeKind.Remove) @@ -437,10 +471,10 @@ public async Task AddOrUpdateStateAsync( return newValue; } - var conditionalResult = await this.TryGetStateFromStateProviderAsync(stateName, cancellationToken); - if (conditionalResult.HasValue) + var (response, httpStatusCode) = await this.TryGetStateFromStateProviderAsync(stateName, cancellationToken); + if (httpStatusCode == 200 && response is not null) { - var newValue = updateValueFactory.Invoke(stateName, conditionalResult.Value.Value); + var newValue = updateValueFactory.Invoke(stateName, response.Value); stateChangeTracker.Add(stateName, StateMetadata.Create(newValue, StateChangeKind.Update, ttl: ttl)); return newValue; @@ -475,7 +509,7 @@ public async Task SaveStateAsync(CancellationToken cancellationToken = default) { var stateMetadata = stateChangeTracker[stateName]; - if (stateMetadata.ChangeKind != StateChangeKind.None) + if (stateMetadata.ChangeKind is not (StateChangeKind.None or StateChangeKind.NotFound)) { stateChangeList.Add( new ActorStateChange(stateName, stateMetadata.Type, stateMetadata.Value, stateMetadata.ChangeKind, stateMetadata.TTLExpireTime)); @@ -500,7 +534,7 @@ public async Task SaveStateAsync(CancellationToken cancellationToken = default) } } - // Remove the states from tracker whcih were marked for removal. + // Remove the states from tracker which were marked for removal. foreach (var stateToRemove in statesToRemove) { stateChangeTracker.Remove(stateToRemove); @@ -561,7 +595,7 @@ private bool IsStateMarkedForRemove(string stateName) return false; } - private Task>> TryGetStateFromStateProviderAsync(string stateName, CancellationToken cancellationToken) + private Task<(ActorStateResponse Response, int HttpStatusCode)> TryGetStateFromStateProviderAsync(string stateName, CancellationToken cancellationToken) { EnsureStateProviderInitialized(); return this.actor.Host.StateProvider.TryLoadStateAsync(this.actorTypeName, this.actor.Id.ToString(), stateName, cancellationToken); @@ -583,15 +617,14 @@ private Dictionary GetContextualStateTracker() { return context.Value.tracker; } - else - { - return defaultTracker; - } + + return defaultTracker; } private sealed class StateMetadata { - private StateMetadata(object value, Type type, StateChangeKind changeKind, DateTimeOffset? ttlExpireTime = null, TimeSpan? ttl = null) + private StateMetadata(object value, Type type, StateChangeKind changeKind, DateTimeOffset? ttlExpireTime = null, + TimeSpan? ttl = null) { this.Value = value; this.Type = type; @@ -601,6 +634,7 @@ private StateMetadata(object value, Type type, StateChangeKind changeKind, DateT { throw new ArgumentException("Cannot specify both TTLExpireTime and TTL"); } + if (ttl.HasValue) { this.TTLExpireTime = DateTimeOffset.UtcNow.Add(ttl.Value); @@ -619,32 +653,22 @@ private StateMetadata(object value, Type type, StateChangeKind changeKind, DateT public DateTimeOffset? TTLExpireTime { get; set; } - public static StateMetadata Create(T value, StateChangeKind changeKind) - { - return new StateMetadata(value, typeof(T), changeKind); - } + public static StateMetadata Create(T value, StateChangeKind changeKind) => new(value, typeof(T), changeKind); - public static StateMetadata Create(T value, StateChangeKind changeKind, DateTimeOffset? ttlExpireTime) - { - return new StateMetadata(value, typeof(T), changeKind, ttlExpireTime: ttlExpireTime); - } + public static StateMetadata Create(T value, StateChangeKind changeKind, DateTimeOffset? ttlExpireTime) => + new(value, typeof(T), changeKind, ttlExpireTime: ttlExpireTime); // Non-generic counterpart to Create for callers (like SyncDefaultTracker) that // already have a boxed value and its runtime Type from an ActorStateChange, rather // than a compile-time T. - public static StateMetadata CreateFromValueAndType(object value, Type type, StateChangeKind changeKind, DateTimeOffset? ttlExpireTime) - { - return new StateMetadata(value, type, changeKind, ttlExpireTime: ttlExpireTime); - } + public static StateMetadata CreateFromValueAndType(object value, Type type, StateChangeKind changeKind, DateTimeOffset? ttlExpireTime) => + new(value, type, changeKind, ttlExpireTime: ttlExpireTime); - public static StateMetadata Create(T value, StateChangeKind changeKind, TimeSpan? ttl) - { - return new StateMetadata(value, typeof(T), changeKind, ttl: ttl); - } + public static StateMetadata Create(T value, StateChangeKind changeKind, TimeSpan? ttl) => + new(value, typeof(T), changeKind, ttl: ttl); - public static StateMetadata CreateForRemove() - { - return new StateMetadata(null, typeof(object), StateChangeKind.Remove); - } + public static StateMetadata CreateNotFound() => new(null, typeof(object), StateChangeKind.NotFound); + + public static StateMetadata CreateForRemove() => new(null, typeof(object), StateChangeKind.Remove); } -} \ No newline at end of file +} diff --git a/src/Dapr.Actors/Runtime/DaprStateProvider.cs b/src/Dapr.Actors/Runtime/DaprStateProvider.cs index bddbaf564..188d1f250 100644 --- a/src/Dapr.Actors/Runtime/DaprStateProvider.cs +++ b/src/Dapr.Actors/Runtime/DaprStateProvider.cs @@ -44,12 +44,11 @@ public DaprStateProvider(IDaprInteractor daprInteractor, JsonSerializerOptions j this.daprInteractor = daprInteractor; } - public async Task>> TryLoadStateAsync(string actorType, string actorId, string stateName, CancellationToken cancellationToken = default) + public async Task<(ActorStateResponse Response, int HttpStatusCode)> TryLoadStateAsync(string actorType, string actorId, string stateName, CancellationToken cancellationToken = default) { - var result = new ConditionalValue>(false, default); var response = await this.daprInteractor.GetStateAsync(actorType, actorId, stateName, cancellationToken); - if (response.Value.Length != 0 && (!response.TTLExpireTime.HasValue || response.TTLExpireTime.Value > DateTimeOffset.UtcNow)) + if (response.HttpStatusCode == 200 && response.Value.Length != 0 && (!response.TTLExpireTime.HasValue || response.TTLExpireTime.Value > DateTimeOffset.UtcNow)) { T typedResult; @@ -64,10 +63,10 @@ public async Task>> TryLoadStateAsync( typedResult = JsonSerializer.Deserialize(response.Value, jsonSerializerOptions); } - result = new ConditionalValue>(true, new ActorStateResponse(typedResult, response.TTLExpireTime)); + return (new ActorStateResponse(typedResult, response.TTLExpireTime, response.HttpStatusCode), response.HttpStatusCode); } - return result; + return (null, response.HttpStatusCode); } public async Task ContainsStateAsync(string actorType, string actorId, string stateName, CancellationToken cancellationToken = default) @@ -102,12 +101,12 @@ private async Task DoStateChangesTransactionallyAsync(string actorType, string a ] */ using var stream = new MemoryStream(); - using var writer = new Utf8JsonWriter(stream); + await using var writer = new Utf8JsonWriter(stream); writer.WriteStartArray(); foreach (var stateChange in stateChanges) { writer.WriteStartObject(); - var operation = this.GetDaprStateOperation(stateChange.ChangeKind); + var operation = GetDaprStateOperation(stateChange.ChangeKind); writer.WriteString("operation", operation); // write the requestProperty @@ -142,8 +141,6 @@ private async Task DoStateChangesTransactionallyAsync(string actorType, string a writer.WriteEndObject(); } - break; - default: break; } @@ -153,12 +150,12 @@ private async Task DoStateChangesTransactionallyAsync(string actorType, string a writer.WriteEndArray(); - await writer.FlushAsync(); + await writer.FlushAsync(cancellationToken); var content = Encoding.UTF8.GetString(stream.ToArray()); await this.daprInteractor.SaveStateTransactionallyAsync(actorType, actorId, content, cancellationToken); } - private string GetDaprStateOperation(StateChangeKind changeKind) + private static string GetDaprStateOperation(StateChangeKind changeKind) { var operation = string.Empty; @@ -171,10 +168,8 @@ private string GetDaprStateOperation(StateChangeKind changeKind) case StateChangeKind.Update: operation = "upsert"; break; - default: - break; } return operation; } -} \ No newline at end of file +} diff --git a/src/Dapr.Actors/Runtime/StateChangeKind.cs b/src/Dapr.Actors/Runtime/StateChangeKind.cs index 3ad79893b..aa3e8afd6 100644 --- a/src/Dapr.Actors/Runtime/StateChangeKind.cs +++ b/src/Dapr.Actors/Runtime/StateChangeKind.cs @@ -37,4 +37,9 @@ public enum StateChangeKind /// The state needs to be removed. /// Remove = 3, -} \ No newline at end of file + + /// + /// There is no state. + /// + NotFound = 4, +} diff --git a/test/Dapr.Actors.Test/ActorStateManagerTest.cs b/test/Dapr.Actors.Test/ActorStateManagerTest.cs index 24b3d288a..e95901dd3 100644 --- a/test/Dapr.Actors.Test/ActorStateManagerTest.cs +++ b/test/Dapr.Actors.Test/ActorStateManagerTest.cs @@ -922,4 +922,337 @@ public async Task SetStateContext_TwoContextsHaveIndependentTrackers() Assert.Equal("from-ctx2", ctx2Result); Assert.False(ctx1Saw); } + + [Theory] + [InlineData(200)] + [InlineData(204)] + [InlineData(400)] + [InlineData(500)] + public async Task TryGetStateAsync_CachesMissingStateAndDoesNotRequeryStateProvider(int statusCode) + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("", null, statusCode))); + + for (var i = 0; i < 100; i++) + { + Assert.False((await mngr.TryGetStateAsync("missing", token)).HasValue); + } + + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "missing", It.IsAny()), + Times.Once); + } + + [Fact] + public async Task TryGetStateAsync_CachesLoadedStateAndDoesNotRequeryStateProvider() + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), "key", It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("\"value\"", null))); + + for (var i = 0; i < 100; i++) + { + var result = await mngr.TryGetStateAsync("key", token); + Assert.True(result.HasValue); + Assert.Equal("value", result.Value); + } + + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "key", It.IsAny()), + Times.Once); + } + + [Fact] + public async Task MissingStateCacheIsScopedByStateName() + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("", null, 204))); + + Assert.False((await mngr.TryGetStateAsync("missing-1", token)).HasValue); + Assert.False((await mngr.TryGetStateAsync("missing-2", token)).HasValue); + Assert.False((await mngr.TryGetStateAsync("missing-1", token)).HasValue); + Assert.False((await mngr.TryGetStateAsync("missing-2", token)).HasValue); + + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "missing-1", It.IsAny()), + Times.Once); + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "missing-2", It.IsAny()), + Times.Once); + } + + [Fact] + public async Task ClearCacheAsync_InvalidatesMissingStateCache() + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + var stateExists = false; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), "key", It.IsAny())) + .ReturnsAsync(() => stateExists + ? new ActorStateResponse("\"value\"", null) + : new ActorStateResponse("", null, 204)); + + Assert.False((await mngr.TryGetStateAsync("key", token)).HasValue); + stateExists = true; + Assert.False((await mngr.TryGetStateAsync("key", token)).HasValue); + + await mngr.ClearCacheAsync(token); + + Assert.Equal("value", await mngr.GetStateAsync("key", token)); + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "key", It.IsAny()), + Times.Exactly(2)); + } + + [Fact] + public async Task UnloadStateAsync_InvalidatesMissingStateCache() + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + var stateExists = false; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), "key", It.IsAny())) + .ReturnsAsync(() => stateExists + ? new ActorStateResponse("\"value\"", null) + : new ActorStateResponse("", null, 204)); + + Assert.False((await mngr.TryGetStateAsync("key", token)).HasValue); + stateExists = true; + Assert.False((await mngr.TryGetStateAsync("key", token)).HasValue); + + await mngr.UnloadStateAsync("key", cancellationToken: token); + + Assert.Equal("value", await mngr.GetStateAsync("key", token)); + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "key", It.IsAny()), + Times.Exactly(2)); + } + + [Fact] + public async Task MissingStateIsNotPersistedOnSave() + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("", null, 204))); + + Assert.False((await mngr.TryGetStateAsync("missing", token)).HasValue); + await mngr.SaveStateAsync(token); + + interactor.Verify( + d => d.SaveStateTransactionallyAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny()), + Times.Never); + } + + [Fact] + public async Task MissingStateIsReportedAsAbsentByContainsAndRemove() + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("", null, 204))); + + Assert.False((await mngr.TryGetStateAsync("missing", token)).HasValue); + + Assert.False(await mngr.ContainsStateAsync("missing", token)); + Assert.False(await mngr.TryRemoveStateAsync("missing", token)); + await Assert.ThrowsAsync(() => mngr.GetStateAsync("missing", token)); + } + + [Theory] + [InlineData(true)] + [InlineData(false)] + public async Task MissingStateCanBeAddedAndIsPersistedAsAnAdd(bool useTtl) + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("", null, 204))); + + string capturedContent = null; + interactor + .Setup(d => d.SaveStateTransactionallyAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Callback((_, _, content, _) => capturedContent = content) + .Returns(Task.CompletedTask); + + Assert.False((await mngr.TryGetStateAsync("missing", token)).HasValue); + + if (useTtl) + { + Assert.True(await mngr.TryAddStateAsync("missing", "value1", TimeSpan.FromMinutes(5), token)); + } + else + { + Assert.True(await mngr.TryAddStateAsync("missing", "value1", token)); + } + + Assert.Equal("value1", await mngr.GetStateAsync("missing", token)); + Assert.True(await mngr.ContainsStateAsync("missing", token)); + + await mngr.SaveStateAsync(token); + Assert.Contains("\"operation\":\"upsert\"", capturedContent); + Assert.Contains("\"value\":\"value1\"", capturedContent); + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "missing", It.IsAny()), + Times.Once); + } + + [Theory] + [InlineData(true)] + [InlineData(false)] + public async Task MissingStateCanBeSetAndIsPersisted(bool useTtl) + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("", null, 204))); + + string capturedContent = null; + interactor + .Setup(d => d.SaveStateTransactionallyAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Callback((_, _, content, _) => capturedContent = content) + .Returns(Task.CompletedTask); + + Assert.False((await mngr.TryGetStateAsync("missing", token)).HasValue); + + if (useTtl) + { + await mngr.SetStateAsync("missing", "value1", TimeSpan.FromMinutes(5), token); + } + else + { + await mngr.SetStateAsync("missing", "value1", token); + } + + Assert.Equal("value1", await mngr.GetStateAsync("missing", token)); + + await mngr.SaveStateAsync(token); + Assert.Contains("\"operation\":\"upsert\"", capturedContent); + Assert.Contains("\"value\":\"value1\"", capturedContent); + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "missing", It.IsAny()), + Times.Once); + } + + [Theory] + [InlineData(true)] + [InlineData(false)] + public async Task MissingStateIsAddedByAddOrUpdateWithoutInvokingUpdateFactory(bool useTtl) + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("", null, 204))); + + Assert.False((await mngr.TryGetStateAsync("missing", token)).HasValue); + + var updateFactoryInvoked = false; + string result; + if (useTtl) + { + result = await mngr.AddOrUpdateStateAsync("missing", "added", (_, existing) => + { + updateFactoryInvoked = true; + return existing; + }, TimeSpan.FromMinutes(5), token); + } + else + { + result = await mngr.AddOrUpdateStateAsync("missing", "added", (_, existing) => + { + updateFactoryInvoked = true; + return existing; + }, token); + } + + Assert.Equal("added", result); + Assert.False(updateFactoryInvoked); + Assert.Equal("added", await mngr.GetStateAsync("missing", token)); + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "missing", It.IsAny()), + Times.Once); + } + + [Theory] + [InlineData(true)] + [InlineData(false)] + public async Task MissingStateIsAddedByGetOrAdd(bool useTtl) + { + var interactor = new Mock(); + var host = ActorHost.CreateForTest(); + host.StateProvider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var mngr = new ActorStateManager(new TestActor(host)); + var token = CancellationToken.None; + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("", null, 204))); + + Assert.False((await mngr.TryGetStateAsync("missing", token)).HasValue); + + var value = useTtl + ? await mngr.GetOrAddStateAsync("missing", "added", TimeSpan.FromMinutes(5), token) + : await mngr.GetOrAddStateAsync("missing", "added", token); + + Assert.Equal("added", value); + Assert.Equal("added", await mngr.GetStateAsync("missing", token)); + Assert.True(await mngr.ContainsStateAsync("missing", token)); + interactor.Verify( + d => d.GetStateAsync(It.IsAny(), It.IsAny(), "missing", It.IsAny()), + Times.Once); + } } diff --git a/test/Dapr.Actors.Test/DaprStateProviderTest.cs b/test/Dapr.Actors.Test/DaprStateProviderTest.cs index 370f284c4..4fe081cf2 100644 --- a/test/Dapr.Actors.Test/DaprStateProviderTest.cs +++ b/test/Dapr.Actors.Test/DaprStateProviderTest.cs @@ -99,32 +99,49 @@ public async Task TryLoadStateAsync() .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) .Returns(Task.FromResult(new ActorStateResponse("", null))); var resp = await provider.TryLoadStateAsync("actorType", "actorId", "key", token); - Assert.False(resp.HasValue); + Assert.Null(resp.Response); interactor .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) .Returns(Task.FromResult(new ActorStateResponse("\"value\"", null))); resp = await provider.TryLoadStateAsync("actorType", "actorId", "key", token); - Assert.True(resp.HasValue); - Assert.Equal("value", resp.Value.Value); - Assert.False(resp.Value.TTLExpireTime.HasValue); + Assert.NotNull(resp.Response); + Assert.Equal(200, resp.HttpStatusCode); + Assert.Equal("value", resp.Response.Value); + Assert.False(resp.Response.TTLExpireTime.HasValue); var ttl = DateTime.UtcNow.AddSeconds(1); interactor .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) .Returns(Task.FromResult(new ActorStateResponse("\"value\"", ttl))); resp = await provider.TryLoadStateAsync("actorType", "actorId", "key", token); - Assert.True(resp.HasValue); - Assert.Equal("value", resp.Value.Value); - Assert.True(resp.Value.TTLExpireTime.HasValue); - Assert.Equal(ttl, resp.Value.TTLExpireTime.Value); + Assert.NotNull(resp.Response); + Assert.Equal("value", resp.Response.Value); + Assert.True(resp.Response.TTLExpireTime.HasValue); + Assert.Equal(ttl, resp.Response.TTLExpireTime.Value); ttl = DateTime.UtcNow.AddSeconds(-1); interactor .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) .Returns(Task.FromResult(new ActorStateResponse("\"value\"", ttl))); resp = await provider.TryLoadStateAsync("actorType", "actorId", "key", token); - Assert.False(resp.HasValue); + Assert.Null(resp.Response); + } + + [Fact] + public async Task TryLoadStateAsync_ReturnsNoValueForNonSuccessStatusCode() + { + var interactor = new Mock(); + var provider = new DaprStateProvider(interactor.Object, new JsonSerializerOptions()); + var token = new CancellationToken(); + + interactor + .Setup(d => d.GetStateAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())) + .Returns(Task.FromResult(new ActorStateResponse("\"value\"", null, 204))); + + var resp = await provider.TryLoadStateAsync("actorType", "actorId", "key", token); + Assert.Null(resp.Response); + Assert.Equal(204, resp.HttpStatusCode); } [Fact] @@ -195,6 +212,6 @@ public async Task TryLoadStateAsync_ReturnsFalseWhenTTLExpireTimeIsExactlyNow() .Returns(Task.FromResult(new ActorStateResponse("\"value\"", ttl))); var resp = await provider.TryLoadStateAsync("actorType", "actorId", "key", token); - Assert.False(resp.HasValue); + Assert.Null(resp.Response); } } \ No newline at end of file