diff --git a/CHANGELOG.md b/CHANGELOG.md index 245ab431..ed986416 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -358,6 +358,16 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/) one-line type change; `AddCaching` now throws at startup if a leftover `ISerializerProxy` registration is present, rather than ignoring it and silently falling back to JSON. +- **Multi-key commands are checked against the Redis Cluster slot map before they run.** `GetAsync(CacheKey[])`, `GetCacheEntriesAsync`, the multi-key `SetAsync` and `RemoveAsync(CacheKey[])` each reach Redis as one command on one node, so a batch whose keys span slots is answered with an error, which the caches log and report as a miss: migrating an app onto a cluster turned a working batch read into a cache that never hits, with nothing thrown to say so. Each of those paths now compares the slots first and throws `CrossSlotKeysException`, naming the two keys that disagree and how to fix it. A non-clustered server maps every key to the same slot, so nothing changes off a cluster, and a disconnected cache skips the check and keeps answering from its own disconnected branch. +- **`ShardKeyEnabled` leaves a key that already carries a hash tag alone.** The shard strategy used + to wrap every key in braces, so a caller who had placed its own `{tag}` in the key ended up with + `app:s:{authz_{orgid}_groups_1}` — Redis hashes from the first `{` to the next `}`, so the tag was + the accidental `authz_{orgid` rather than the org the caller asked for. A key with a valid hash tag + (non-empty content between the first `{` and the next `}`) is now prefixed and left as-is, so the + caller's tag alone picks the slot and multi-key commands land where the caller intended. Keys with + no braces, or with braces that do not form a valid tag, are wrapped exactly as before. This + relocates existing entries only for apps that run `ShardKeyEnabled: true` *and* put braces in + their keys. ### Removed diff --git a/docs/how-to/telemetry-and-strategies.md b/docs/how-to/telemetry-and-strategies.md index 7e6566e2..5810a1e7 100644 --- a/docs/how-to/telemetry-and-strategies.md +++ b/docs/how-to/telemetry-and-strategies.md @@ -329,7 +329,7 @@ When `NotifyShardedPubSub` is `true` on `RedisStreamsTopicOptions`, the stream k ### When to reach for these seams -Most consumers never override any of these seams. The built-in `DefaultRedisKeyStrategyFactory` covers `ICache` and `IHashCache` and handles shard-key routing automatically when `CacheOptions.ShardKeyEnabled` is `true`. Reasons to override: +Most consumers never override any of these seams. The built-in `DefaultRedisKeyStrategyFactory` covers `ICache` and `IHashCache` and handles shard-key routing automatically when `CacheOptions.ShardKeyEnabled` is `true`: the key is wrapped in a `{...}` hash tag so the slot follows the key alone, and a key that already carries a valid hash tag is left as it is. Reasons to override: - **Tenant/region shard routing** — you want Redis keys to encode tenant or region info as a Redis Cluster hash tag (`{tag}`) so related keys cluster to the same slot and cross-slot multi-key commands work correctly. - **Legacy key layout compatibility** — you are sharing a Redis instance with an existing system that has a fixed key schema you must match. Overriding the factory lets you produce keys that match the legacy layout without changing the cache call sites. diff --git a/docs/reference/interfaces.md b/docs/reference/interfaces.md index 83656be2..a700cd07 100644 --- a/docs/reference/interfaces.md +++ b/docs/reference/interfaces.md @@ -213,6 +213,8 @@ Two reasons the interface is the strict surface. An implementation cannot silent The multi-key `GetOrAddAsync` overloads shown above pair each key with an opaque caller state (`TState`) and are **default interface methods** — each forwards to the shared `BatchGetOrAdd.RunAsync` machinery, so existing `ICache` implementations keep compiling without adding them. The generator is invoked at most once, only with the states of the entries that missed every cache layer, never with keys; results come back keyed by state, one entry per distinct requested state in first-occurrence order. Cache operations de-duplicate by `CacheKey`; results de-duplicate by state — when two states share a key the generator is asked once and the value is reported under both. Their parameter shape matches the single-key `GetOrAddAsync` exactly — `expiration` and `policy` required, `token` optional. `CacheExtensions` carries no token-positional forwarder for them: those extensions are pre-`CachePolicy` back-compat sugar, and this API predates nothing. There is no key-only convenience overload; a caller whose keys are their own identity pairs each key with itself (`TState = CacheKey`). +Every multi-key method on the Redis-backed caches runs as one command against one node: `GetAsync(CacheKey[])` is an `MGET`, `RemoveAsync(CacheKey[])` one `DEL`, and `GetCacheEntriesAsync` and the multi-key `SetAsync` a single `MULTI`/`EXEC`. On Redis Cluster that requires every key in the batch to hash to one slot, so the caches check the slots first and throw `CrossSlotKeysException` when they differ, naming the two keys that disagree. Redis itself answers a cross-slot command with an error, which the caches log and report as a miss, so without the check a batch spanning slots would read as a cache that never hits. Give the keys a shared hash tag (`app:s:{org1}:groups_1`) so they hash together, or split the batch into one call per tag. A non-clustered server maps every key to the same slot, so the check never fires there, and a disconnected cache skips it and answers from its own disconnected branch. + `TryAddAsync` writes only if the key is absent — StackExchange.Redis `When.NotExists` (`SET … NX`) on Redis-backed caches, in one atomic command with the TTL applied by the same write. Four points decide whether it fits your problem: - **`false` is deliberately ambiguous.** It means "you did not create this key" — either it already existed, or the write could not be completed (backing store disconnected, write threw, or the value was a `null`/`default` that the cache has no way to represent). This is fail-closed by design: a caller treating `true` as "I own this key" is never wrongly told it won. The ambiguity is not recoverable: a serialization or command failure also returns `false`, with [`IConnectionState.IsConnected`](#iconnectionstate) still reporting healthy, and [`IDistributedLock.TryAcquireAsync`](#idistributedlock) conflates backend-unavailable with already-held in the same way. Design the `false` branch so that not proceeding is safe; if the two readings must be handled differently, the caller needs a primitive with a richer result than a `bool`. diff --git a/docs/reference/settings.md b/docs/reference/settings.md index 6410ed47..a8039fff 100644 --- a/docs/reference/settings.md +++ b/docs/reference/settings.md @@ -20,7 +20,7 @@ Every binding-visible property on every shipped options class, with shipped defa | `Enabled` | `bool` | `true` | App-wide | Master on/off switch for the caching subsystem. | | `TelemetryEnabled` | `bool` | `true` | App-wide | Gates the `ICachingTelemetryProvider` seam; set to `false` to silence all cache metrics. | | `BroadcastEnabled` | `bool` | `true` | App-wide | Gates the `ITopicFactory` wiring; set to `false` to disable all invalidation broadcasts. | -| `ShardKeyEnabled` | `bool` | `false` | App-wide | Enable for Redis Cluster deployments that span multiple shards. | +| `ShardKeyEnabled` | `bool` | `false` | App-wide | Enable for Redis Cluster deployments that span multiple shards. Wraps the cache key in a `{...}` hash tag so only the key, not `AppShortName` or the keyspace, picks the slot. A key that already carries a valid hash tag is left as it is, so a caller that placed its own `{tag}` keeps the slot it chose. | | `AuditEnabled` | `bool` | `true` | App-wide | Log writes whose serialized size exceeds `LargeValueThreshold` bytes. | | `DefaultCache` | `string` | `"InMemoryRedis"` | App-wide | Provider name resolved when no explicit provider is requested; values: `InMemory`, `Redis`, `InMemoryRedis`. | | `DefaultTopic` | `string` | `"RedisStreams"` | App-wide | Topic provider used when no explicit topic is requested; values: `RedisStreams`, `RedisPubSub`. | diff --git a/src/UiPath.Caching/Broadcast/Redis/StreamSuffixShardedChannelStrategy.cs b/src/UiPath.Caching/Broadcast/Redis/StreamSuffixShardedChannelStrategy.cs index 19b6e62b..80c3b27c 100644 --- a/src/UiPath.Caching/Broadcast/Redis/StreamSuffixShardedChannelStrategy.cs +++ b/src/UiPath.Caching/Broadcast/Redis/StreamSuffixShardedChannelStrategy.cs @@ -21,11 +21,11 @@ public RedisChannel GetRedisChannel(TopicKey topicKey) private static string ResolveChannelBase(string streamKey) { - if (HasValidHashTag(streamKey)) + if (RedisHashTag.HasValidTag(streamKey)) { return streamKey; } - if (ContainsNoBraces(streamKey)) + if (RedisHashTag.ContainsNoBraces(streamKey)) { return "{" + streamKey + "}"; } @@ -34,18 +34,4 @@ private static string ResolveChannelBase(string streamKey) "The stream key must either contain a valid hash tag (non-empty content between '{' and '}', e.g. 'app:st:{topic}') " + "or contain no '{' or '}' characters at all."); } - - private static bool HasValidHashTag(string key) - { - var open = key.IndexOf('{'); - if (open < 0) - { - return false; - } - var close = key.IndexOf('}', open + 1); - return close > open + 1; - } - - private static bool ContainsNoBraces(string key) => - key.IndexOf('{') < 0 && key.IndexOf('}') < 0; } diff --git a/src/UiPath.Caching/PublicAPI.Unshipped.txt b/src/UiPath.Caching/PublicAPI.Unshipped.txt index ba6bfeb5..d4b15c0e 100644 --- a/src/UiPath.Caching/PublicAPI.Unshipped.txt +++ b/src/UiPath.Caching/PublicAPI.Unshipped.txt @@ -153,3 +153,7 @@ UiPath.Caching.CacheEventPublisher.CacheRefreshedAsync(UiPath.Caching.ICacheEntr UiPath.Caching.CacheEventPublisher.CacheRemovedAsync(UiPath.Caching.ICacheEntryOptions! options, System.Type? entryType) -> System.Threading.Tasks.ValueTask UiPath.Caching.CacheEventPublisher.CacheSetAsync(UiPath.Caching.ICacheEntryOptions! options, System.Type? entryType) -> System.Threading.Tasks.ValueTask UiPath.Caching.CacheEventPublisher.MetadataUpdatedAsync(UiPath.Caching.ICacheEntryOptions! options, System.Type? entryType) -> System.Threading.Tasks.ValueTask +UiPath.Caching.Redis.CrossSlotKeysException +UiPath.Caching.Redis.CrossSlotKeysException.CrossSlotKeysException() -> void +UiPath.Caching.Redis.CrossSlotKeysException.CrossSlotKeysException(string? message) -> void +UiPath.Caching.Redis.CrossSlotKeysException.CrossSlotKeysException(string? message, System.Exception? innerException) -> void diff --git a/src/UiPath.Caching/Redis/CrossSlotKeysException.cs b/src/UiPath.Caching/Redis/CrossSlotKeysException.cs new file mode 100644 index 00000000..dfecc545 --- /dev/null +++ b/src/UiPath.Caching/Redis/CrossSlotKeysException.cs @@ -0,0 +1,16 @@ +namespace UiPath.Caching.Redis; + +public class CrossSlotKeysException : Exception +{ + public CrossSlotKeysException(string? message) : base(message) + { + } + + public CrossSlotKeysException() + { + } + + public CrossSlotKeysException(string? message, Exception? innerException) : base(message, innerException) + { + } +} diff --git a/src/UiPath.Caching/Redis/RedisCache.cs b/src/UiPath.Caching/Redis/RedisCache.cs index 74c11782..1085fa7d 100644 --- a/src/UiPath.Caching/Redis/RedisCache.cs +++ b/src/UiPath.Caching/Redis/RedisCache.cs @@ -205,7 +205,7 @@ public ValueTask RemoveAsync(CacheKey cacheKey, CancellationToken token public ValueTask RemoveAsync(CacheKey[] cacheKey, CancellationToken token = default) { - return RemoveAsync(cacheKey.Select(k => ToRedisKey(k, token)).ToArray(), token); + return RemoveAsync(cacheKey, cacheKey.Select(k => ToRedisKey(k, token)).ToArray(), token); } public ValueTask SetAsync(CacheKey cacheKey, T? value, CachePolicy? policy, CancellationToken token = default) => @@ -543,7 +543,8 @@ private async ValueTask SetInternalAsync(KeyValuePair[] k { bool ret = default; token.ThrowIfCancellationRequested(); - + var redisKeys = keyValues.Select(kv => ToRedisKey(kv.Key, token)).ToArray(); + if (!IsConnected) { return false; @@ -552,12 +553,13 @@ private async ValueTask SetInternalAsync(KeyValuePair[] k var operation = StartOperation(nameof(SetAsync)); try { + ThrowIfCrossSlot(Array.ConvertAll(keyValues, kv => kv.Key), redisKeys, typeof(T), nameof(SetAsync)); var transaction = Database.CreateTransaction(asyncState: null); - foreach (var keyValue in keyValues) + for (var i = 0; i < keyValues.Length; i++) { - var redisKey = ToRedisKey(keyValue.Key, token); - var value = keyValue.Value; + var redisKey = redisKeys[i]; + var value = keyValues[i].Value; if (IsDefault(value)) { if (_cacheNullValues && expiration > TimeSpan.Zero) @@ -584,6 +586,12 @@ private async ValueTask SetInternalAsync(KeyValuePair[] k operation.Stop(); } + catch (CrossSlotKeysException) + { + // A batch the caller has to fix, not a Redis failure to report as a miss. + operation.Stop(); + throw; + } catch (Exception ex) { operation.Stop(); @@ -625,13 +633,14 @@ private async ValueTask RemoveAsync(RedisKey redisKey, CancellationToke return ret; } - private async ValueTask RemoveAsync(RedisKey[] redisKey, CancellationToken token) + private async ValueTask RemoveAsync(CacheKey[] cacheKeys, RedisKey[] redisKey, CancellationToken token) { bool ret = default; token.ThrowIfCancellationRequested(); var operation = StartOperation(); try { + ThrowIfCrossSlot(cacheKeys, redisKey, typeof(T), nameof(RemoveAsync)); var response = await _write.ExecuteAsync(async token => { token.ThrowIfCancellationRequested(); @@ -640,6 +649,12 @@ private async ValueTask RemoveAsync(RedisKey[] redisKey, CancellationTo operation.Stop(); ret = response > -1; } + catch (CrossSlotKeysException) + { + // A batch the caller has to fix, not a Redis failure to report as a miss. + operation.Stop(); + throw; + } catch (Exception ex) { operation.Stop(); @@ -707,6 +722,7 @@ private async ValueTask RemoveAsync(RedisKey[] redisKey, CancellationTo var reads = InitReads(redisKeys); try { + ThrowIfCrossSlot(keys, redisKeys, typeof(T), nameof(GetAsync)); var values = await _read.ExecuteAsync(async token => { token.ThrowIfCancellationRequested(); @@ -729,6 +745,12 @@ private async ValueTask RemoveAsync(RedisKey[] redisKey, CancellationTo } operation.Stop(); } + catch (CrossSlotKeysException) + { + // A batch the caller has to fix, not a Redis failure to report as a miss. + operation.Stop(); + throw; + } catch (Exception ex) { operation.Stop(); @@ -829,6 +851,7 @@ private async ValueTask RemoveAsync(RedisKey[] redisKey, CancellationTo var reads = InitReads(redisKeys); try { + ThrowIfCrossSlot(keys, redisKeys, typeof(T), nameof(GetCacheEntriesAsync)); var transaction = Database.CreateTransaction(); var mgetTask = transaction.StringGetAsync(redisKeys, CommandFlags.PreferReplica).ConfigureAwait(false); var (expireTimeTasks, ttlTasks) = StartExpirationFetches(transaction, redisKeys); @@ -865,6 +888,12 @@ private async ValueTask RemoveAsync(RedisKey[] redisKey, CancellationTo } operation.Stop(); } + catch (CrossSlotKeysException) + { + // A batch the caller has to fix, not a Redis failure to report as a miss. + operation.Stop(); + throw; + } catch (Exception ex) { operation.Stop(); diff --git a/src/UiPath.Caching/Redis/RedisCacheBase.cs b/src/UiPath.Caching/Redis/RedisCacheBase.cs index 56028052..f479be17 100644 --- a/src/UiPath.Caching/Redis/RedisCacheBase.cs +++ b/src/UiPath.Caching/Redis/RedisCacheBase.cs @@ -118,6 +118,36 @@ public event EventHandler? OnReconnected public bool IsConnected => _connectionState.IsConnected; + /// + /// Redis answers a cross-slot command with an error the caches log and report as a miss, so a batch + /// spanning slots would read as a cache that never hits. Call it inside the command's own try, since + /// reading a slot resolves the connection, and let back out past + /// the catch that turns a Redis failure into a miss. + /// + private protected void ThrowIfCrossSlot(CacheKey[] cacheKeys, RedisKey[] redisKeys, Type? valueType, [CallerMemberName] string? operation = null) + { + if (redisKeys.Length < 2) + { + return; + } + + var multiplexer = Database.Multiplexer; + var slot = multiplexer.GetHashSlot(redisKeys[0]); + for (var i = 1; i < redisKeys.Length; i++) + { + if (multiplexer.GetHashSlot(redisKeys[i]) == slot) + { + continue; + } + throw new CrossSlotKeysException( + $"{operation} was given {redisKeys.Length} keys that Redis Cluster maps to different slots, " + + $"starting with '{Logged(cacheKeys[0], redisKeys[0], valueType)}' and '{Logged(cacheKeys[i], redisKeys[i], valueType)}'. " + + "A multi-key command runs on one node: give the keys a shared hash tag (non-empty content " + + "between '{' and '}', e.g. 'app:s:{org1}:groups_1') so they hash together, or split the batch " + + "into one call per group of keys that already share a tag."); + } + } + protected IDatabase Database => _redis.Database; public void Dispose() diff --git a/src/UiPath.Caching/Redis/RedisHashTag.cs b/src/UiPath.Caching/Redis/RedisHashTag.cs new file mode 100644 index 00000000..55626287 --- /dev/null +++ b/src/UiPath.Caching/Redis/RedisHashTag.cs @@ -0,0 +1,19 @@ +namespace UiPath.Caching.Redis; + +internal static class RedisHashTag +{ + /// Non-empty content between the first '{' and the next '}' is what Redis Cluster hashes. + public static bool HasValidTag(string key) + { + var open = key.IndexOf('{'); + if (open < 0) + { + return false; + } + var close = key.IndexOf('}', open + 1); + return close > open + 1; + } + + public static bool ContainsNoBraces(string key) => + key.IndexOf('{') < 0 && key.IndexOf('}') < 0; +} diff --git a/src/UiPath.Caching/Redis/ShardPrefixRedisKeyStrategy.cs b/src/UiPath.Caching/Redis/ShardPrefixRedisKeyStrategy.cs index 38cb192b..82a8bc82 100644 --- a/src/UiPath.Caching/Redis/ShardPrefixRedisKeyStrategy.cs +++ b/src/UiPath.Caching/Redis/ShardPrefixRedisKeyStrategy.cs @@ -9,6 +9,9 @@ public ShardPrefixRedisKeyStrategy(string prefix, char separator) : base(prefix, { } + // Wrapping a key that already carries a tag would hash the tag together with what precedes it. public override RedisKey GetRedisKey(CacheKey key) => - string.Join(Separator, Prefix, string.Format(CultureInfo.InvariantCulture, ShardFormat, key)); + RedisHashTag.HasValidTag(key.Name) + ? base.GetRedisKey(key) + : string.Join(Separator, Prefix, string.Format(CultureInfo.InvariantCulture, ShardFormat, key)); } diff --git a/tests/UiPath.Caching.Tests/Redis/RedisCacheTests.cs b/tests/UiPath.Caching.Tests/Redis/RedisCacheTests.cs index 69afdfbe..b9cb3de3 100644 --- a/tests/UiPath.Caching.Tests/Redis/RedisCacheTests.cs +++ b/tests/UiPath.Caching.Tests/Redis/RedisCacheTests.cs @@ -84,6 +84,59 @@ public async Task Multi_get_works_as_expected() actualValue.Should().BeEquivalentTo(new KeyValuePair[] { new(_cacheKey, expectedValue), new(_multiKey, expectedValue) }); } + [Fact] + public async Task Multi_get_throws_when_keys_span_slots() + { + GiveKeysDifferentSlots(); + var act = () => Sut.GetAsync(new CacheKey[] { _cacheKey, _multiKey }, policy: null, token: testContextAccessor.Current.CancellationToken).AsTask(); + (await act.Should().ThrowAsync()).And.Message.Should().Contain("GetAsync"); + await _database.DidNotReceive().StringGetAsync(Arg.Any(), Arg.Any()); + } + + [Fact] + public async Task Multi_get_cache_entries_throws_when_keys_span_slots() + { + GiveKeysDifferentSlots(); + var act = () => Sut.GetCacheEntriesAsync(new CacheKey[] { _cacheKey, _multiKey }, policy: null, token: testContextAccessor.Current.CancellationToken).AsTask(); + await act.Should().ThrowAsync(); + _database.DidNotReceive().CreateTransaction(); + } + + [Fact] + public async Task Multi_remove_throws_when_keys_span_slots() + { + GiveKeysDifferentSlots(); + var act = () => Sut.RemoveAsync(new CacheKey[] { _cacheKey, _multiKey }, testContextAccessor.Current.CancellationToken).AsTask(); + await act.Should().ThrowAsync(); + await _database.DidNotReceive().KeyDeleteAsync(Arg.Any(), Arg.Any()); + } + + [Fact] + public async Task Multi_set_throws_when_keys_span_slots() + { + GiveKeysDifferentSlots(); + var value = _fixture.Create(); + var act = () => Sut.SetAsync(new KeyValuePair[] { new(_cacheKey, value), new(_multiKey, value) }, _fixture.Create(), policy: null, token: testContextAccessor.Current.CancellationToken).AsTask(); + await act.Should().ThrowAsync(); + _database.DidNotReceive().CreateTransaction(Arg.Any()); + } + + [Fact] + public async Task Multi_get_still_reports_a_miss_when_the_connection_cannot_be_resolved() + { + // The slot check reads the multiplexer, so a Redis that cannot be reached must not turn a + // graceful miss into a thrown connection error. + _connector.Database.Throws(new RedisException("test")); + var actual = await Sut.GetAsync(new CacheKey[] { _cacheKey, _multiKey }, policy: null, token: testContextAccessor.Current.CancellationToken); + actual.Should().BeEquivalentTo(new KeyValuePair[] { new(_cacheKey, null), new(_multiKey, null) }); + } + + private void GiveKeysDifferentSlots() + { + _database.Multiplexer.GetHashSlot(_redisKey).Returns(1); + _database.Multiplexer.GetHashSlot(_redisMultiKey).Returns(2); + } + [Fact] public async Task Multi_get_has_no_redis_exceptions() { @@ -1393,6 +1446,8 @@ public ValueTask InitializeAsync() }; _database = _fixture.Freeze(); + // A non-clustered server answers every key with the same slot; the fixture's auto-values would hand each one a random slot. + _database.Multiplexer.GetHashSlot(Arg.Any()).Returns(0); _transaction = _fixture.Freeze(); _database.CreateTransaction().Returns(_transaction); _serializer = new SystemJsonByteSerializerProxy(); diff --git a/tests/UiPath.Caching.Tests/ShardPrefixRedisKeyStrategyTests.cs b/tests/UiPath.Caching.Tests/ShardPrefixRedisKeyStrategyTests.cs index e5d41108..f7a6f782 100644 --- a/tests/UiPath.Caching.Tests/ShardPrefixRedisKeyStrategyTests.cs +++ b/tests/UiPath.Caching.Tests/ShardPrefixRedisKeyStrategyTests.cs @@ -15,6 +15,10 @@ public class ShardPrefixRedisKeyStrategyTests [Theory] [InlineData("xxx", '$', "bla", "xxx${bla}")] [InlineData("aa", 'B', "ccc", "aab{ccc}")] + [InlineData("app:s", ':', "{org1}:groups_1", "app:s:{org1}:groups_1")] + [InlineData("app:s", ':', "authz_{org1}_groups_1", "app:s:authz_{org1}_groups_1")] + [InlineData("app:s", ':', "bla{}bla", "app:s:{bla{}bla}")] + [InlineData("app:s", ':', "bla{bla", "app:s:{bla{bla}")] public void WorksAsExpected(string prefix, char separator, string key, string expected) { _prefix = prefix;