diff --git a/docs/Streams.md b/docs/Streams.md index 47e82c2b9..0121392d7 100644 --- a/docs/Streams.md +++ b/docs/Streams.md @@ -37,6 +37,21 @@ You also have the option to override the auto-generated message ID by passing yo db.StreamAdd("events_stream", "foo_name", "bar_value", messageId: "0-1", maxLength: 100); ``` +Trimming, the entry ID and the other `XADD` options can be supplied together via `StreamAddOptions`: + +```csharp +// append only if the stream already exists, trimming anything older than the given entry ID +var id = db.StreamAdd("events_stream", "foo_name", "bar_value", new StreamAddOptions +{ + CreateStream = false, // NOMKSTREAM + MinId = "1526919030474-55", // or MaxLength, for MAXLEN + Approximate = true, // the "~" form, which is cheaper +}); +``` + +With `CreateStream = false`, adding to a stream that does not exist returns `RedisValue.Null` and the key is +not created. `MinId` and `Approximate` require server version 6.2 or above. + Idempotent write-at-most-once production === diff --git a/src/StackExchange.Redis/Interfaces/IDatabase.cs b/src/StackExchange.Redis/Interfaces/IDatabase.cs index 038439a0b..141ed8bf4 100644 --- a/src/StackExchange.Redis/Interfaces/IDatabase.cs +++ b/src/StackExchange.Redis/Interfaces/IDatabase.cs @@ -2835,6 +2835,41 @@ IEnumerable SortedSetScan( /// The ID of the newly created message. /// RedisValue StreamAdd(RedisKey key, NameValueEntry[] streamPairs, StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode trimMode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None); + + /// + /// Adds an entry using the specified values to the given stream key. + /// If key does not exist and is set, a new key holding a + /// stream is created. The command returns the ID of the newly created stream entry. + /// + /// The key of the stream. + /// The field name for the stream entry. + /// The value to set in the stream entry. + /// Additional options for this operation, such as trimming and the entry ID. + /// The flags to use for this operation. + /// + /// The ID of the newly created message, or a null value when the key does not exist and + /// is false. + /// + /// +#pragma warning disable RS0027 // additive overload: `options` is required, so existing calls still bind to the overloads above + RedisValue StreamAdd(RedisKey key, RedisValue streamField, RedisValue streamValue, StreamAddOptions options, CommandFlags flags = CommandFlags.None); + + /// + /// Adds an entry using the specified values to the given stream key. + /// If key does not exist and is set, a new key holding a + /// stream is created. The command returns the ID of the newly created stream entry. + /// + /// The key of the stream. + /// The fields and their associated values to set in the stream entry. + /// Additional options for this operation, such as trimming and the entry ID. + /// The flags to use for this operation. + /// + /// The ID of the newly created message, or a null value when the key does not exist and + /// is false. + /// + /// + RedisValue StreamAdd(RedisKey key, NameValueEntry[] streamPairs, StreamAddOptions options, CommandFlags flags = CommandFlags.None); +#pragma warning restore RS0027 #pragma warning restore RS0026 /// diff --git a/src/StackExchange.Redis/Interfaces/IDatabaseAsync.cs b/src/StackExchange.Redis/Interfaces/IDatabaseAsync.cs index f80aa744a..84403ab32 100644 --- a/src/StackExchange.Redis/Interfaces/IDatabaseAsync.cs +++ b/src/StackExchange.Redis/Interfaces/IDatabaseAsync.cs @@ -698,6 +698,14 @@ IAsyncEnumerable SortedSetScanAsync( /// Task StreamAddAsync(RedisKey key, NameValueEntry[] streamPairs, StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode trimMode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None); + + /// +#pragma warning disable RS0027 // additive overload: `options` is required, so existing calls still bind to the overloads above + Task StreamAddAsync(RedisKey key, RedisValue streamField, RedisValue streamValue, StreamAddOptions options, CommandFlags flags = CommandFlags.None); + + /// + Task StreamAddAsync(RedisKey key, NameValueEntry[] streamPairs, StreamAddOptions options, CommandFlags flags = CommandFlags.None); +#pragma warning restore RS0027 #pragma warning restore RS0026 /// diff --git a/src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixed.cs b/src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixed.cs index c92f24d5c..d2b1a815d 100644 --- a/src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixed.cs +++ b/src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixed.cs @@ -670,6 +670,12 @@ public Task StreamAddAsync(RedisKey key, RedisValue streamField, Red public Task StreamAddAsync(RedisKey key, NameValueEntry[] streamPairs, StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) => Inner.StreamAddAsync(ToInner(key), streamPairs, idempotentId, maxLength, useApproximateMaxLength, limit, mode, flags); + public Task StreamAddAsync(RedisKey key, RedisValue streamField, RedisValue streamValue, StreamAddOptions options, CommandFlags flags = CommandFlags.None) => + Inner.StreamAddAsync(ToInner(key), streamField, streamValue, options, flags); + + public Task StreamAddAsync(RedisKey key, NameValueEntry[] streamPairs, StreamAddOptions options, CommandFlags flags = CommandFlags.None) => + Inner.StreamAddAsync(ToInner(key), streamPairs, options, flags); + public Task StreamConfigureAsync(RedisKey key, StreamConfiguration configuration, CommandFlags flags = CommandFlags.None) => Inner.StreamConfigureAsync(ToInner(key), configuration, flags); diff --git a/src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixedDatabase.cs b/src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixedDatabase.cs index cbf7c72fc..1ba10ecca 100644 --- a/src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixedDatabase.cs +++ b/src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixedDatabase.cs @@ -637,6 +637,12 @@ public RedisValue StreamAdd(RedisKey key, RedisValue streamField, RedisValue str public RedisValue StreamAdd(RedisKey key, NameValueEntry[] streamPairs, StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) => Inner.StreamAdd(ToInner(key), streamPairs, idempotentId, maxLength, useApproximateMaxLength, limit, mode, flags); + public RedisValue StreamAdd(RedisKey key, RedisValue streamField, RedisValue streamValue, StreamAddOptions options, CommandFlags flags = CommandFlags.None) => + Inner.StreamAdd(ToInner(key), streamField, streamValue, options, flags); + + public RedisValue StreamAdd(RedisKey key, NameValueEntry[] streamPairs, StreamAddOptions options, CommandFlags flags = CommandFlags.None) => + Inner.StreamAdd(ToInner(key), streamPairs, options, flags); + public void StreamConfigure(RedisKey key, StreamConfiguration configuration, CommandFlags flags = CommandFlags.None) => Inner.StreamConfigure(ToInner(key), configuration, flags); diff --git a/src/StackExchange.Redis/PublicAPI/PublicAPI.Unshipped.txt b/src/StackExchange.Redis/PublicAPI/PublicAPI.Unshipped.txt index 295fc5b1f..3e6df8527 100644 --- a/src/StackExchange.Redis/PublicAPI/PublicAPI.Unshipped.txt +++ b/src/StackExchange.Redis/PublicAPI/PublicAPI.Unshipped.txt @@ -36,3 +36,25 @@ StackExchange.Redis.IServer.ClusterSlots(StackExchange.Redis.CommandFlags flags StackExchange.Redis.IServer.ClusterSlotsAsync(StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> System.Threading.Tasks.Task! [SER007]StackExchange.Redis.RedisErrorKind.UnknownRedirectTarget = 26 -> StackExchange.Redis.RedisErrorKind override StackExchange.Redis.ClusterSlotNode.ToString() -> string! +StackExchange.Redis.IDatabaseAsync.StreamAddAsync(StackExchange.Redis.RedisKey key, StackExchange.Redis.NameValueEntry[]! streamPairs, StackExchange.Redis.StreamAddOptions options, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> System.Threading.Tasks.Task! +StackExchange.Redis.IDatabaseAsync.StreamAddAsync(StackExchange.Redis.RedisKey key, StackExchange.Redis.RedisValue streamField, StackExchange.Redis.RedisValue streamValue, StackExchange.Redis.StreamAddOptions options, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> System.Threading.Tasks.Task! +StackExchange.Redis.IDatabase.StreamAdd(StackExchange.Redis.RedisKey key, StackExchange.Redis.NameValueEntry[]! streamPairs, StackExchange.Redis.StreamAddOptions options, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> StackExchange.Redis.RedisValue +StackExchange.Redis.IDatabase.StreamAdd(StackExchange.Redis.RedisKey key, StackExchange.Redis.RedisValue streamField, StackExchange.Redis.RedisValue streamValue, StackExchange.Redis.StreamAddOptions options, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> StackExchange.Redis.RedisValue +StackExchange.Redis.StreamAddOptions +StackExchange.Redis.StreamAddOptions.Approximate.get -> bool +StackExchange.Redis.StreamAddOptions.Approximate.init -> void +StackExchange.Redis.StreamAddOptions.CreateStream.get -> bool +StackExchange.Redis.StreamAddOptions.CreateStream.init -> void +StackExchange.Redis.StreamAddOptions.IdempotentId.get -> StackExchange.Redis.StreamIdempotentId +StackExchange.Redis.StreamAddOptions.IdempotentId.init -> void +StackExchange.Redis.StreamAddOptions.Limit.get -> long? +StackExchange.Redis.StreamAddOptions.Limit.init -> void +StackExchange.Redis.StreamAddOptions.MaxLength.get -> long? +StackExchange.Redis.StreamAddOptions.MaxLength.init -> void +StackExchange.Redis.StreamAddOptions.MessageId.get -> StackExchange.Redis.RedisValue? +StackExchange.Redis.StreamAddOptions.MessageId.init -> void +StackExchange.Redis.StreamAddOptions.MinId.get -> StackExchange.Redis.RedisValue +StackExchange.Redis.StreamAddOptions.MinId.init -> void +StackExchange.Redis.StreamAddOptions.StreamAddOptions() -> void +StackExchange.Redis.StreamAddOptions.TrimMode.get -> StackExchange.Redis.StreamTrimMode +StackExchange.Redis.StreamAddOptions.TrimMode.init -> void diff --git a/src/StackExchange.Redis/RedisDatabase.cs b/src/StackExchange.Redis/RedisDatabase.cs index 9e4c18ac7..91c99d566 100644 --- a/src/StackExchange.Redis/RedisDatabase.cs +++ b/src/StackExchange.Redis/RedisDatabase.cs @@ -2896,33 +2896,22 @@ public RedisValue StreamAdd(RedisKey key, RedisValue streamField, RedisValue str public RedisValue StreamAdd(RedisKey key, RedisValue streamField, RedisValue streamValue, RedisValue? messageId = null, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) { - var msg = GetStreamAddMessage( - key, - messageId ?? StreamConstants.AutoGeneratedId, - StreamIdempotentId.Empty, - maxLength, - useApproximateMaxLength, - new NameValueEntry(streamField, streamValue), - limit, - mode, - flags); - + var options = LegacyStreamAddOptions(messageId, StreamIdempotentId.Empty, maxLength, useApproximateMaxLength, limit, mode); + var msg = GetStreamAddMessage(key, in options, new NameValueEntry(streamField, streamValue), flags); return ExecuteSync(msg, ResultProcessor.RedisValue); } public RedisValue StreamAdd(RedisKey key, RedisValue streamField, RedisValue streamValue, StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) { - var msg = GetStreamAddMessage( - key, - StreamConstants.AutoGeneratedId, - idempotentId, - maxLength, - useApproximateMaxLength, - new NameValueEntry(streamField, streamValue), - limit, - mode, - flags); + var options = LegacyStreamAddOptions(null, in idempotentId, maxLength, useApproximateMaxLength, limit, mode); + var msg = GetStreamAddMessage(key, in options, new NameValueEntry(streamField, streamValue), flags); + return ExecuteSync(msg, ResultProcessor.RedisValue); + } + public RedisValue StreamAdd(RedisKey key, RedisValue streamField, RedisValue streamValue, StreamAddOptions options, CommandFlags flags = CommandFlags.None) + { + options.ThrowIfInvalid(); + var msg = GetStreamAddMessage(key, in options, new NameValueEntry(streamField, streamValue), flags); return ExecuteSync(msg, ResultProcessor.RedisValue); } @@ -2931,33 +2920,22 @@ public Task StreamAddAsync(RedisKey key, RedisValue streamField, Red public Task StreamAddAsync(RedisKey key, RedisValue streamField, RedisValue streamValue, RedisValue? messageId = null, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) { - var msg = GetStreamAddMessage( - key, - messageId ?? StreamConstants.AutoGeneratedId, - StreamIdempotentId.Empty, - maxLength, - useApproximateMaxLength, - new NameValueEntry(streamField, streamValue), - limit, - mode, - flags); - + var options = LegacyStreamAddOptions(messageId, StreamIdempotentId.Empty, maxLength, useApproximateMaxLength, limit, mode); + var msg = GetStreamAddMessage(key, in options, new NameValueEntry(streamField, streamValue), flags); return ExecuteAsync(msg, ResultProcessor.RedisValue); } public Task StreamAddAsync(RedisKey key, RedisValue streamField, RedisValue streamValue, StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) { - var msg = GetStreamAddMessage( - key, - StreamConstants.AutoGeneratedId, - idempotentId, - maxLength, - useApproximateMaxLength, - new NameValueEntry(streamField, streamValue), - limit, - mode, - flags); + var options = LegacyStreamAddOptions(null, in idempotentId, maxLength, useApproximateMaxLength, limit, mode); + var msg = GetStreamAddMessage(key, in options, new NameValueEntry(streamField, streamValue), flags); + return ExecuteAsync(msg, ResultProcessor.RedisValue); + } + public Task StreamAddAsync(RedisKey key, RedisValue streamField, RedisValue streamValue, StreamAddOptions options, CommandFlags flags = CommandFlags.None) + { + options.ThrowIfInvalid(); + var msg = GetStreamAddMessage(key, in options, new NameValueEntry(streamField, streamValue), flags); return ExecuteAsync(msg, ResultProcessor.RedisValue); } @@ -2966,33 +2944,22 @@ public RedisValue StreamAdd(RedisKey key, NameValueEntry[] streamPairs, RedisVal public RedisValue StreamAdd(RedisKey key, NameValueEntry[] streamPairs, RedisValue? messageId = null, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) { - var msg = GetStreamAddMessage( - key, - messageId ?? StreamConstants.AutoGeneratedId, - StreamIdempotentId.Empty, - maxLength, - useApproximateMaxLength, - streamPairs, - limit, - mode, - flags); - + var options = LegacyStreamAddOptions(messageId, StreamIdempotentId.Empty, maxLength, useApproximateMaxLength, limit, mode); + var msg = GetStreamAddMessage(key, in options, streamPairs, flags); return ExecuteSync(msg, ResultProcessor.RedisValue); } public RedisValue StreamAdd(RedisKey key, NameValueEntry[] streamPairs, StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) { - var msg = GetStreamAddMessage( - key, - StreamConstants.AutoGeneratedId, - idempotentId, - maxLength, - useApproximateMaxLength, - streamPairs, - limit, - mode, - flags); + var options = LegacyStreamAddOptions(null, in idempotentId, maxLength, useApproximateMaxLength, limit, mode); + var msg = GetStreamAddMessage(key, in options, streamPairs, flags); + return ExecuteSync(msg, ResultProcessor.RedisValue); + } + public RedisValue StreamAdd(RedisKey key, NameValueEntry[] streamPairs, StreamAddOptions options, CommandFlags flags = CommandFlags.None) + { + options.ThrowIfInvalid(); + var msg = GetStreamAddMessage(key, in options, streamPairs, flags); return ExecuteSync(msg, ResultProcessor.RedisValue); } @@ -3001,33 +2968,22 @@ public Task StreamAddAsync(RedisKey key, NameValueEntry[] streamPair public Task StreamAddAsync(RedisKey key, NameValueEntry[] streamPairs, RedisValue? messageId = null, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) { - var msg = GetStreamAddMessage( - key, - messageId ?? StreamConstants.AutoGeneratedId, - StreamIdempotentId.Empty, - maxLength, - useApproximateMaxLength, - streamPairs, - limit, - mode, - flags); - + var options = LegacyStreamAddOptions(messageId, StreamIdempotentId.Empty, maxLength, useApproximateMaxLength, limit, mode); + var msg = GetStreamAddMessage(key, in options, streamPairs, flags); return ExecuteAsync(msg, ResultProcessor.RedisValue); } public Task StreamAddAsync(RedisKey key, NameValueEntry[] streamPairs, StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StreamTrimMode mode = StreamTrimMode.KeepReferences, CommandFlags flags = CommandFlags.None) { - var msg = GetStreamAddMessage( - key, - StreamConstants.AutoGeneratedId, - idempotentId, - maxLength, - useApproximateMaxLength, - streamPairs, - limit, - mode, - flags); + var options = LegacyStreamAddOptions(null, in idempotentId, maxLength, useApproximateMaxLength, limit, mode); + var msg = GetStreamAddMessage(key, in options, streamPairs, flags); + return ExecuteAsync(msg, ResultProcessor.RedisValue); + } + public Task StreamAddAsync(RedisKey key, NameValueEntry[] streamPairs, StreamAddOptions options, CommandFlags flags = CommandFlags.None) + { + options.ThrowIfInvalid(); + var msg = GetStreamAddMessage(key, in options, streamPairs, flags); return ExecuteAsync(msg, ResultProcessor.RedisValue); } @@ -4909,53 +4865,92 @@ private Message GetStreamAcknowledgeAndDeleteMessage(RedisKey key, RedisValue gr return Message.Create(Database, flags, RedisCommand.XACKDEL, key, values); } - internal Message GetStreamAddMessage(in RedisKey key, RedisValue messageId, in StreamIdempotentId idempotentId, long? maxLength, bool useApproximateMaxLength, NameValueEntry streamPair, long? limit, StreamTrimMode mode, CommandFlags flags) - { - // Calculate the correct number of arguments: - // 3 array elements for Entry ID & NameValueEntry.Name & NameValueEntry.Value. - // 2 elements if using MAXLEN (keyword & value), otherwise 0. - // 1 element if using Approximate Length (~), otherwise 0. - var totalLength = 3 + (maxLength.HasValue ? 2 : 0) - + idempotentId.ArgCount - + (maxLength.HasValue && useApproximateMaxLength ? 1 : 0) - + (limit.HasValue ? 2 : 0) - + (mode != StreamTrimMode.KeepReferences ? 1 : 0); + /// + /// Maps the legacy positional trim/id arguments onto . + /// + /// + /// Deliberately does *not* call : the shipped overloads + /// have always passed questionable combinations (LIMIT without a threshold, say) through to the + /// server, and that behaviour is preserved; only the options-based overloads validate up-front. + /// + private static StreamAddOptions LegacyStreamAddOptions(RedisValue? messageId, in StreamIdempotentId idempotentId, long? maxLength, bool useApproximateMaxLength, long? limit, StreamTrimMode mode) + => new() + { + MessageId = messageId, + IdempotentId = idempotentId, + MaxLength = maxLength, + Approximate = useApproximateMaxLength, + Limit = limit, + TrimMode = mode, + }; - var values = new RedisValue[totalLength]; - var offset = 0; + /// + /// The number of arguments written by . + /// + private static int GetStreamAddPrefixLength(in StreamAddOptions options) + => (options.CreateStream ? 0 : 1) // NOMKSTREAM + + (options.HasThreshold ? 2 : 0) // MAXLEN|MINID + + (options.HasThreshold && options.Approximate ? 1 : 0) // ~ + + (options.Limit.HasValue ? 2 : 0) // LIMIT N + + (options.TrimMode == StreamTrimMode.KeepReferences ? 0 : 1) // relevant trim-mode keyword + + options.IdempotentId.ArgCount // IDMP / IDMPAUTO + + 1; // the stream entry ID + + /// + /// Writes everything in XADD between the key and the field/value pairs, i.e. + /// [NOMKSTREAM] [MAXLEN|MINID [~] threshold] [LIMIT n] [KEEPREF|DELREF|ACKED] [IDMP...] <*|id>. + /// + private static void WriteStreamAddPrefix(in StreamAddOptions options, RedisValue[] values, ref int offset) + { + if (!options.CreateStream) + { + values[offset++] = StreamConstants.NoMkStream; + } - if (maxLength.HasValue) + if (options.HasThreshold) { - values[offset++] = StreamConstants.MaxLen; + var byMaxLength = options.MaxLength.HasValue; + values[offset++] = byMaxLength ? StreamConstants.MaxLen : StreamConstants.MinId; - if (useApproximateMaxLength) + if (options.Approximate) { values[offset++] = StreamConstants.ApproximateMaxLen; } - values[offset++] = maxLength.Value; + values[offset++] = byMaxLength ? options.MaxLength.GetValueOrDefault() : options.MinId; } - if (limit.HasValue) + if (options.Limit.HasValue) { values[offset++] = RedisLiterals.LIMIT; - values[offset++] = limit.Value; + values[offset++] = options.Limit.GetValueOrDefault(); } - if (mode != StreamTrimMode.KeepReferences) + if (options.TrimMode != StreamTrimMode.KeepReferences) { - values[offset++] = StreamConstants.GetMode(mode); + values[offset++] = StreamConstants.GetMode(options.TrimMode); } - idempotentId.WriteTo(values, ref offset); + options.IdempotentId.WriteTo(values, ref offset); + + values[offset++] = options.EntryId; + } + + internal Message GetStreamAddMessage(in RedisKey key, in StreamAddOptions options, NameValueEntry streamPair, CommandFlags flags) + { + // 2 array elements for NameValueEntry.Name & NameValueEntry.Value, after the shared prefix + var totalLength = 2 + GetStreamAddPrefixLength(in options); + + var values = new RedisValue[totalLength]; + var offset = 0; - values[offset++] = messageId; + WriteStreamAddPrefix(in options, values, ref offset); values[offset++] = streamPair.Name; values[offset++] = streamPair.Value; Debug.Assert(offset == totalLength); - return Message.Create(Database, GetStreamAddCategory(flags, messageId, in idempotentId), RedisCommand.XADD, key, values); + return Message.Create(Database, GetStreamAddCategory(flags, in options), RedisCommand.XADD, key, values); } /// @@ -4969,8 +4964,8 @@ internal Message GetStreamAddMessage(in RedisKey key, RedisValue messageId, in S /// appends a second entry (5-0, then 5-1) rather than being rejected. Testing only against the bare /// * would therefore let a double-append through under the default policy. /// - private static CommandFlags GetStreamAddCategory(CommandFlags flags, in RedisValue messageId, in StreamIdempotentId idempotentId) - => (idempotentId.ArgCount != 0 || !IsServerAssignedId(in messageId)) + private static CommandFlags GetStreamAddCategory(CommandFlags flags, in StreamAddOptions options) + => (options.IdempotentId.ArgCount != 0 || !IsServerAssignedId(options.EntryId)) ? flags.WithCategory(CommandFlags.CommandRetryWriteChecked) : flags; @@ -4985,54 +4980,23 @@ internal static bool IsServerAssignedId(in RedisValue messageId) /// /// Gets message for . /// - private Message GetStreamAddMessage(in RedisKey key, RedisValue entryId, in StreamIdempotentId idempotentId, long? maxLength, bool useApproximateMaxLength, NameValueEntry[] streamPairs, long? limit, StreamTrimMode mode, CommandFlags flags) + internal Message GetStreamAddMessage(in RedisKey key, in StreamAddOptions options, NameValueEntry[] streamPairs, CommandFlags flags) { if (streamPairs == null) throw new ArgumentNullException(nameof(streamPairs)); if (streamPairs.Length == 0) throw new ArgumentOutOfRangeException(nameof(streamPairs), "streamPairs must contain at least one item."); - if (maxLength.HasValue && maxLength <= 0) + if (options.MaxLength.HasValue && options.MaxLength <= 0) { - throw new ArgumentOutOfRangeException(nameof(maxLength), "maxLength must be greater than 0."); + throw new ArgumentOutOfRangeException(nameof(options), "maxLength must be greater than 0."); } var totalLength = (streamPairs.Length * 2) // Room for the name/value pairs - + 1 // The stream entry ID - + idempotentId.ArgCount - + (maxLength.HasValue ? 2 : 0) // MAXLEN N - + (maxLength.HasValue && useApproximateMaxLength ? 1 : 0) // ~ - + (mode == StreamTrimMode.KeepReferences ? 0 : 1) // relevant trim-mode keyword - + (limit.HasValue ? 2 : 0); // LIMIT N + + GetStreamAddPrefixLength(in options); var values = new RedisValue[totalLength]; - var offset = 0; - if (maxLength.HasValue) - { - values[offset++] = StreamConstants.MaxLen; - - if (useApproximateMaxLength) - { - values[offset++] = StreamConstants.ApproximateMaxLen; - } - - values[offset++] = maxLength.Value; - } - - if (limit.HasValue) - { - values[offset++] = RedisLiterals.LIMIT; - values[offset++] = limit.Value; - } - - if (mode != StreamTrimMode.KeepReferences) - { - values[offset++] = StreamConstants.GetMode(mode); - } - - idempotentId.WriteTo(values, ref offset); - - values[offset++] = entryId; + WriteStreamAddPrefix(in options, values, ref offset); for (var i = 0; i < streamPairs.Length; i++) { @@ -5041,7 +5005,7 @@ private Message GetStreamAddMessage(in RedisKey key, RedisValue entryId, in Stre } Debug.Assert(offset == totalLength); - return Message.Create(Database, GetStreamAddCategory(flags, entryId, in idempotentId), RedisCommand.XADD, key, values); + return Message.Create(Database, GetStreamAddCategory(flags, in options), RedisCommand.XADD, key, values); } internal Message GetStreamAutoClaimMessage(RedisKey key, RedisValue consumerGroup, RedisValue assignToConsumer, long minIdleTimeInMs, RedisValue startAtId, int? count, bool idsOnly, CommandFlags flags) diff --git a/src/StackExchange.Redis/StreamAddOptions.cs b/src/StackExchange.Redis/StreamAddOptions.cs new file mode 100644 index 000000000..7ccec8898 --- /dev/null +++ b/src/StackExchange.Redis/StreamAddOptions.cs @@ -0,0 +1,105 @@ +using System; + +namespace StackExchange.Redis; + +/// +/// Additional options for adding an entry to a stream; see . +/// +/// +/// The default value of this type means "append a new entry, letting the server assign the id, creating the +/// stream if it does not already exist, without trimming" - i.e. a plain XADD key * .... +/// +public readonly struct StreamAddOptions +{ + // stored inverted, so that `default` means CreateStream=true; XADD creates the stream unless NOMKSTREAM + // is specified, and the default of this type must match the default of the command + private readonly bool _noMkStream; + + /// + /// The ID to assign to the stream entry; defaults to an auto-generated ID ("*"). + /// + /// Mutually exclusive with . + public RedisValue? MessageId { get; init; } + + /// + /// The idempotent producer (pid) and optionally id (iid) to use for this entry; the server assigns the + /// entry ID in this mode. See for more information of the idempotent API. + /// + /// Mutually exclusive with . + public StreamIdempotentId IdempotentId { get; init; } + + /// + /// Trim the stream to this maximum length (MAXLEN) after appending. + /// + /// Mutually exclusive with . + public long? MaxLength { get; init; } + + /// + /// Trim entries with an ID lower than this (MINID) after appending. + /// + /// Mutually exclusive with . + public RedisValue MinId { get; init; } + + /// + /// If true, the "~" argument is used to allow the stream to exceed the requested threshold by a small + /// number. This improves performance when removing messages. Ignored when no threshold is specified. + /// + public bool Approximate { get; init; } + + /// + /// Specifies the maximal count of entries that will be evicted; requires + /// and a threshold. + /// + public long? Limit { get; init; } + + /// + /// Determines how stream trimming should be performed. + /// + public StreamTrimMode TrimMode { get; init; } + + /// + /// Whether to create the stream when the key does not exist; when false, NOMKSTREAM is used and + /// the command returns a null value instead of creating the stream. Defaults to true. + /// + public bool CreateStream + { + get => !_noMkStream; + init => _noMkStream = !value; + } + + /// + /// Whether a trim threshold (MAXLEN or MINID) has been specified. + /// + internal bool HasThreshold => MaxLength.HasValue || MinId.HasValue; + + /// + /// The entry ID to send; the idempotent forms always let the server assign the ID. + /// + internal RedisValue EntryId => MessageId ?? StreamConstants.AutoGeneratedId; + + internal void ThrowIfInvalid() + { + if (MaxLength.HasValue && MinId.HasValue) + { + Throw($"{nameof(MaxLength)} and {nameof(MinId)} are mutually exclusive."); + } + if (MessageId.HasValue && IdempotentId.ArgCount != 0) + { + Throw($"{nameof(MessageId)} and {nameof(IdempotentId)} are mutually exclusive; the server assigns the entry ID when producing idempotently."); + } + if (Limit.HasValue) + { + // both of these are rejected by the server; fail earlier, and more clearly + if (!HasThreshold) + { + Throw($"{nameof(Limit)} requires {nameof(MaxLength)} or {nameof(MinId)}."); + } + if (!Approximate) + { + Throw($"{nameof(Limit)} requires {nameof(Approximate)}."); + } + } + + static void Throw(string message) => throw new ArgumentException(message, "options"); + } +} diff --git a/src/StackExchange.Redis/StreamConstants.cs b/src/StackExchange.Redis/StreamConstants.cs index b63489e01..ba98ef2da 100644 --- a/src/StackExchange.Redis/StreamConstants.cs +++ b/src/StackExchange.Redis/StreamConstants.cs @@ -63,6 +63,8 @@ internal static class StreamConstants internal static readonly RedisValue MkStream = RedisValue.FromRaw("MKSTREAM"u8); + internal static readonly RedisValue NoMkStream = RedisValue.FromRaw("NOMKSTREAM"u8); + internal static readonly RedisValue Stream = RedisValue.FromRaw("STREAM"u8); private static readonly RedisValue KeepRef = RedisValue.FromRaw("KEEPREF"u8), DelRef = RedisValue.FromRaw("DELREF"u8), Acked = RedisValue.FromRaw("ACKED"u8); diff --git a/tests/StackExchange.Redis.Tests/CommandRetryCategoryUnitTests.cs b/tests/StackExchange.Redis.Tests/CommandRetryCategoryUnitTests.cs index 836a00a7b..6ab6341ec 100644 --- a/tests/StackExchange.Redis.Tests/CommandRetryCategoryUnitTests.cs +++ b/tests/StackExchange.Redis.Tests/CommandRetryCategoryUnitTests.cs @@ -1,4 +1,4 @@ -using System; +using System; using System.Threading.Tasks; using Xunit; @@ -110,6 +110,10 @@ public async Task SortedSetAdd_IncrementAccumulates() AssertCategory(Checked, db.GetSortedSetIncrementMessage(key, member, 1.0, ValueCondition.NotExists, CommandFlags.None), "ZADD NX INCR"); } + + private static StreamAddOptions Options(RedisValue messageId, in StreamIdempotentId idempotentId) => + new() { MessageId = messageId, IdempotentId = idempotentId }; + [Fact] public async Task StreamAdd_ExplicitAndIdempotentIdsAreReplaySafe() { @@ -121,27 +125,27 @@ public async Task StreamAdd_ExplicitAndIdempotentIdsAreReplaySafe() // "*" lets the server pick the id, so a replay appends a second entry AssertCategory( Accumulating, - db.GetStreamAddMessage(key, "*", in noId, null, false, pair, null, StreamTrimMode.KeepReferences, CommandFlags.None), + db.GetStreamAddMessage(key, Options("*", in noId), pair, CommandFlags.None), "XADD *"); // an explicit id is rejected second time round ("equal or smaller") AssertCategory( Checked, - db.GetStreamAddMessage(key, "5-5", in noId, null, false, pair, null, StreamTrimMode.KeepReferences, CommandFlags.None), + db.GetStreamAddMessage(key, Options("5-5", in noId), pair, CommandFlags.None), "XADD with explicit id"); // IDMP producer id: the server deduplicates var idmp = new StreamIdempotentId("producer", "item-1"); AssertCategory( Checked, - db.GetStreamAddMessage(key, "*", in idmp, null, false, pair, null, StreamTrimMode.KeepReferences, CommandFlags.None), + db.GetStreamAddMessage(key, Options("*", in idmp), pair, CommandFlags.None), "XADD IDMP"); // IDMPAUTO producer: same, with the id derived from the entry content var idmpAuto = new StreamIdempotentId("producer"); AssertCategory( Checked, - db.GetStreamAddMessage(key, "*", in idmpAuto, null, false, pair, null, StreamTrimMode.KeepReferences, CommandFlags.None), + db.GetStreamAddMessage(key, Options("*", in idmpAuto), pair, CommandFlags.None), "XADD IDMPAUTO"); } @@ -159,7 +163,7 @@ public async Task StreamAdd_PartialAutoIdStillAccumulates() var pair = new NameValueEntry("f", "v"); var noId = default(StreamIdempotentId); - Message Add(RedisValue id) => db.GetStreamAddMessage(key, id, in noId, null, false, pair, null, StreamTrimMode.KeepReferences, CommandFlags.None); + Message Add(RedisValue id) => db.GetStreamAddMessage(key, Options(id, in noId), pair, CommandFlags.None); // anything the server completes accumulates... AssertCategory(Accumulating, Add("*"), "XADD *"); diff --git a/tests/StackExchange.Redis.Tests/KeyPrefixedDatabaseTests.cs b/tests/StackExchange.Redis.Tests/KeyPrefixedDatabaseTests.cs index a93c425dc..fa7e1fba2 100644 --- a/tests/StackExchange.Redis.Tests/KeyPrefixedDatabaseTests.cs +++ b/tests/StackExchange.Redis.Tests/KeyPrefixedDatabaseTests.cs @@ -1777,6 +1777,23 @@ public void StreamAdd_WithTrimMode_2() mock.Received().StreamAdd("prefix:key", fields, "*", 1000, false, 100, StreamTrimMode.KeepReferences, CommandFlags.None); } + [Fact] + public void StreamAdd_WithOptions_1() + { + var options = new StreamAddOptions { MaxLength = 1000, CreateStream = false }; + prefixed.StreamAdd("key", "field", "value", options, CommandFlags.None); + mock.Received().StreamAdd("prefix:key", "field", "value", options, CommandFlags.None); + } + + [Fact] + public void StreamAdd_WithOptions_2() + { + var fields = new NameValueEntry[] { new NameValueEntry("field", "value") }; + var options = new StreamAddOptions { MinId = "5-5", CreateStream = false }; + prefixed.StreamAdd("key", fields, options, CommandFlags.None); + mock.Received().StreamAdd("prefix:key", fields, options, CommandFlags.None); + } + [Fact] public void StreamTrim_WithMode() { diff --git a/tests/StackExchange.Redis.Tests/KeyPrefixedTests.cs b/tests/StackExchange.Redis.Tests/KeyPrefixedTests.cs index 8ab933deb..5df1b380d 100644 --- a/tests/StackExchange.Redis.Tests/KeyPrefixedTests.cs +++ b/tests/StackExchange.Redis.Tests/KeyPrefixedTests.cs @@ -1705,6 +1705,23 @@ public async Task StreamAddAsync_WithTrimMode_2() await mock.Received().StreamAddAsync("prefix:key", fields, "*", 1000, false, 100, StreamTrimMode.KeepReferences, CommandFlags.None); } + [Fact] + public async Task StreamAddAsync_WithOptions_1() + { + var options = new StreamAddOptions { MaxLength = 1000, CreateStream = false }; + await prefixed.StreamAddAsync("key", "field", "value", options, CommandFlags.None); + await mock.Received().StreamAddAsync("prefix:key", "field", "value", options, CommandFlags.None); + } + + [Fact] + public async Task StreamAddAsync_WithOptions_2() + { + var fields = new NameValueEntry[] { new NameValueEntry("field", "value") }; + var options = new StreamAddOptions { MinId = "5-5", CreateStream = false }; + await prefixed.StreamAddAsync("key", fields, options, CommandFlags.None); + await mock.Received().StreamAddAsync("prefix:key", fields, options, CommandFlags.None); + } + [Fact] public async Task StreamTrimAsync_WithMode() { diff --git a/tests/StackExchange.Redis.Tests/RoundTripUnitTests/StreamAddRoundTrip.cs b/tests/StackExchange.Redis.Tests/RoundTripUnitTests/StreamAddRoundTrip.cs new file mode 100644 index 000000000..13969734e --- /dev/null +++ b/tests/StackExchange.Redis.Tests/RoundTripUnitTests/StreamAddRoundTrip.cs @@ -0,0 +1,129 @@ +using System; +using System.Threading.Tasks; +using Xunit; + +namespace StackExchange.Redis.Tests.RoundTripUnitTests; + +/// +/// Pins the exact XADD argument order, which the server is strict about: +/// [NOMKSTREAM] [MAXLEN|MINID [~] threshold] [LIMIT n] [KEEPREF|DELREF|ACKED] [IDMP...] <*|id>. +/// +public class StreamAddRoundTrip(ITestOutputHelper log) +{ + private const string Reply = "$3\r\n1-0\r\n"; + + [Fact] + public Task NoOptions() => AssertPairAsync( + default, + "*5\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$1\r\n*\r\n$5\r\nfield\r\n$5\r\nvalue\r\n"); + + [Fact] + public Task NoMkStream() => AssertPairAsync( + new() { CreateStream = false }, + "*6\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$10\r\nNOMKSTREAM\r\n$1\r\n*\r\n$5\r\nfield\r\n$5\r\nvalue\r\n"); + + /// NOMKSTREAM comes before the trim options, not after them. + [Fact] + public Task NoMkStreamPrecedesMaxLen() => AssertPairAsync( + new() { CreateStream = false, MaxLength = 10, Approximate = true, Limit = 5 }, + "*11\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$10\r\nNOMKSTREAM\r\n$6\r\nMAXLEN\r\n$1\r\n~\r\n$2\r\n10\r\n$5\r\nLIMIT\r\n$1\r\n5\r\n$1\r\n*\r\n$5\r\nfield\r\n$5\r\nvalue\r\n"); + + [Fact] + public Task MaxLenOnly() => AssertPairAsync( + new() { MaxLength = 10 }, + "*7\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$6\r\nMAXLEN\r\n$2\r\n10\r\n$1\r\n*\r\n$5\r\nfield\r\n$5\r\nvalue\r\n"); + + [Fact] + public Task MinIdExact() => AssertPairAsync( + new() { MinId = "1526919030474-55" }, + "*7\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$5\r\nMINID\r\n$16\r\n1526919030474-55\r\n$1\r\n*\r\n$5\r\nfield\r\n$5\r\nvalue\r\n"); + + [Fact] + public Task MinIdApproximateWithLimitAndTrimMode() => AssertPairAsync( + new() { MinId = "5-5", Approximate = true, Limit = 3, TrimMode = StreamTrimMode.DeleteReferences }, + "*11\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$5\r\nMINID\r\n$1\r\n~\r\n$3\r\n5-5\r\n$5\r\nLIMIT\r\n$1\r\n3\r\n$6\r\nDELREF\r\n$1\r\n*\r\n$5\r\nfield\r\n$5\r\nvalue\r\n"); + + /// An explicit id replaces the "*", and everything else keeps its place. + [Fact] + public Task ExplicitMessageId() => AssertPairAsync( + new() { MessageId = "5-5", CreateStream = false }, + "*6\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$10\r\nNOMKSTREAM\r\n$3\r\n5-5\r\n$5\r\nfield\r\n$5\r\nvalue\r\n"); + + /// The idempotency arguments sit between the trim options and the entry id. + [Fact] + public Task NoMkStreamWithIdempotentId() => AssertPairAsync( + new() { CreateStream = false, IdempotentId = new("producer", "item-1") }, + "*9\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$10\r\nNOMKSTREAM\r\n$4\r\nIDMP\r\n$8\r\nproducer\r\n$6\r\nitem-1\r\n$1\r\n*\r\n$5\r\nfield\r\n$5\r\nvalue\r\n"); + + [Fact] + public Task NoMkStreamWithAutoIdempotentIdAndMaxLen() => AssertPairAsync( + new() { CreateStream = false, IdempotentId = new("producer"), MaxLength = 10 }, + "*10\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$10\r\nNOMKSTREAM\r\n$6\r\nMAXLEN\r\n$2\r\n10\r\n$8\r\nIDMPAUTO\r\n$8\r\nproducer\r\n$1\r\n*\r\n$5\r\nfield\r\n$5\r\nvalue\r\n"); + + /// The multi-pair builder is a separate code path, and needs its own arithmetic checked. + [Fact] + public async Task PairsArray() + { + var db = new RedisDatabase(null!, 0, null); + NameValueEntry[] pairs = [new("f1", "v1"), new("f2", "v2")]; + var options = new StreamAddOptions { CreateStream = false, MinId = "5-5", Approximate = true }; + var message = db.GetStreamAddMessage("stream", in options, pairs, CommandFlags.None); + + var result = await TestConnection.ExecuteAsync( + message, + ResultProcessor.RedisValue, + "*11\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$10\r\nNOMKSTREAM\r\n$5\r\nMINID\r\n$1\r\n~\r\n$3\r\n5-5\r\n$1\r\n*\r\n$2\r\nf1\r\n$2\r\nv1\r\n$2\r\nf2\r\n$2\r\nv2\r\n", + Reply, + log: log); + Assert.Equal("1-0", result); + } + + /// With NOMKSTREAM and no stream, the server replies nil - which must surface as a null value. + [Fact] + public async Task NilReplyIsNullValue() + { + var db = new RedisDatabase(null!, 0, null); + var options = new StreamAddOptions { CreateStream = false }; + var message = db.GetStreamAddMessage("stream", in options, new NameValueEntry("field", "value"), CommandFlags.None); + + var result = await TestConnection.ExecuteAsync( + message, + ResultProcessor.RedisValue, + "*6\r\n$4\r\nXADD\r\n$6\r\nstream\r\n$10\r\nNOMKSTREAM\r\n$1\r\n*\r\n$5\r\nfield\r\n$5\r\nvalue\r\n", + "$-1\r\n", + log: log); + Assert.True(result.IsNull); + } + + [Theory] + [InlineData(nameof(StreamAddOptions.MaxLength))] + [InlineData(nameof(StreamAddOptions.MessageId))] + [InlineData("LimitWithoutThreshold")] + [InlineData("LimitWithoutApproximate")] + public void InvalidCombinationsAreRejected(string scenario) + { + var db = new RedisDatabase(null!, 0, null); + var options = scenario switch + { + nameof(StreamAddOptions.MaxLength) => new StreamAddOptions { MaxLength = 10, MinId = "5-5" }, + nameof(StreamAddOptions.MessageId) => new StreamAddOptions { MessageId = "5-5", IdempotentId = new("producer") }, + "LimitWithoutThreshold" => new StreamAddOptions { Limit = 5, Approximate = true }, + _ => new StreamAddOptions { MaxLength = 10, Limit = 5 }, + }; + + // the validation is on the public entry-points, not the message builder: the shipped positional + // overloads have always passed odd-but-legal-looking combinations to the server, and still do + var ex = Assert.Throws(() => db.StreamAdd("stream", "field", "value", options)); + log.WriteLine(ex.Message); + Assert.Equal("options", ex.ParamName); + } + + private async Task AssertPairAsync(StreamAddOptions options, string requestResp) + { + var db = new RedisDatabase(null!, 0, null); + var message = db.GetStreamAddMessage("stream", in options, new NameValueEntry("field", "value"), CommandFlags.None); + + var result = await TestConnection.ExecuteAsync(message, ResultProcessor.RedisValue, requestResp, Reply, log: log); + Assert.Equal("1-0", result); + } +} diff --git a/tests/StackExchange.Redis.Tests/StreamTests.cs b/tests/StackExchange.Redis.Tests/StreamTests.cs index 514d21763..0a62d86ec 100644 --- a/tests/StackExchange.Redis.Tests/StreamTests.cs +++ b/tests/StackExchange.Redis.Tests/StreamTests.cs @@ -81,6 +81,77 @@ public async Task StreamAddWithManualId() Assert.Equal(id, messageId); } + [Theory] + [InlineData(false, false)] + [InlineData(false, true)] + [InlineData(true, false)] + [InlineData(true, true)] + public async Task StreamAddCreateStreamFalse(bool pairs, bool useAsync) + { + await using var conn = Create(require: RedisFeatures.v6_2_0); + var db = conn.GetDatabase(); + var key = Me() + $":{pairs}:{useAsync}"; + await db.KeyDeleteAsync(key); + + async Task Add(bool createStream) + { + var options = new StreamAddOptions { CreateStream = createStream }; + if (pairs) + { + NameValueEntry[] fields = [new("field1", "value1"), new("field2", "value2")]; + return useAsync + ? await db.StreamAddAsync(key, fields, options) + : db.StreamAdd(key, fields, options); + } + + return useAsync + ? await db.StreamAddAsync(key, "field", "value", options) + : db.StreamAdd(key, "field", "value", options); + } + + // no stream, and we declined to create one: nothing happens, and we are told so + Assert.Equal(RedisValue.Null, await Add(createStream: false)); + Assert.False(await db.KeyExistsAsync(key)); + + // ...but once the stream exists, the same call appends as normal + Assert.NotEqual(RedisValue.Null, await Add(createStream: true)); + Assert.NotEqual(RedisValue.Null, await Add(createStream: false)); + Assert.Equal(2, await db.StreamLengthAsync(key)); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task StreamAddTrimsByMinId(bool approximate) + { + await using var conn = Create(require: RedisFeatures.v6_2_0); + var db = conn.GetDatabase(); + var key = Me() + $":{approximate}"; + await db.KeyDeleteAsync(key); + + for (var i = 1; i <= 5; i++) + { + db.StreamAdd(key, "f", i, new StreamAddOptions { MessageId = $"{i}-1" }); + } + Assert.Equal(5, await db.StreamLengthAsync(key)); + + // exact MINID must drop everything below the threshold; the approximate form is + // free to keep more, so only assert what the server guarantees in each mode + db.StreamAdd(key, "f", 6, new StreamAddOptions { MessageId = "6-1", MinId = "4-1", Approximate = approximate }); + + var entries = await db.StreamRangeAsync(key); + Assert.Contains(entries, x => x.Id == "6-1"); + if (approximate) + { + Assert.True(entries.Length <= 6); + } + else + { + Assert.Equal(3, entries.Length); // 4-1, 5-1, 6-1 + Assert.DoesNotContain(entries, x => x.Id == "3-1"); + } + } + [Theory] [InlineData(false, false, false)] [InlineData(false, false, true)]