Skip to content

Commit 64999bd

Browse files
authored
Merge branch 'main' into marc/zinter_count
2 parents 4a02c65 + 2dab656 commit 64999bd

12 files changed

Lines changed: 304 additions & 0 deletions

File tree

docs/ReleaseNotes.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ Current package versions:
99
## Unreleased
1010

1111
- Detect server-mode correctly on Valkey 8+ instances ([#3050 by @wipiano](https://github.com/StackExchange/StackExchange.Redis/pull/3050))
12+
- Add Redis 8.8 stream negative acknowledgements (`XNACK`) ([#3058 by @mgravell](https://github.com/StackExchange/StackExchange.Redis/pull/3058))
1213
- Update experimental `GCRA` APIs and wire protocol terminology from "requests" to "tokens", to match server change ([#3051 by @mgravell](https://github.com/StackExchange/StackExchange.Redis/pull/3051))
1314
- Add experimental `Aggregate.Count` support for sorted-set combination operations against Redis 8.8 ([#3059 by @mgravell](https://github.com/StackExchange/StackExchange.Redis/pull/3059))
1415

src/StackExchange.Redis/Enums/RedisCommand.cs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -240,6 +240,7 @@ internal enum RedisCommand
240240
XGROUP,
241241
XINFO,
242242
XLEN,
243+
XNACK,
243244
XPENDING,
244245
XRANGE,
245246
XREAD,
@@ -561,6 +562,7 @@ internal static bool IsPrimaryOnly(this RedisCommand command)
561562
case RedisCommand.XDEL:
562563
case RedisCommand.XDELEX:
563564
case RedisCommand.XGROUP:
565+
case RedisCommand.XNACK:
564566
case RedisCommand.XREADGROUP:
565567
case RedisCommand.XTRIM:
566568
return false;
Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,26 @@
1+
using System.Diagnostics.CodeAnalysis;
2+
using RESPite;
3+
4+
namespace StackExchange.Redis;
5+
6+
/// <summary>
7+
/// Determines how a stream message is negatively acknowledged back to the consumer group.
8+
/// </summary>
9+
[Experimental(Experiments.Server_8_8, UrlFormat = Experiments.UrlFormat)]
10+
public enum StreamNackMode
11+
{
12+
/// <summary>
13+
/// Release the message without counting it as an additional failure.
14+
/// </summary>
15+
Silent = 0,
16+
17+
/// <summary>
18+
/// Release the message and treat it as a normal failed delivery.
19+
/// </summary>
20+
Fail = 1,
21+
22+
/// <summary>
23+
/// Release the message and mark it as a terminal failure.
24+
/// </summary>
25+
Fatal = 2,
26+
}

src/StackExchange.Redis/Interfaces/IDatabase.cs

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2614,6 +2614,36 @@ IEnumerable<SortedSetEntry> SortedSetScan(
26142614
StreamTrimResult[] StreamAcknowledgeAndDelete(RedisKey key, RedisValue groupName, StreamTrimMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None);
26152615
#pragma warning restore RS0026
26162616

2617+
/// <summary>
2618+
/// Allow the consumer to release a pending message back to the group without marking it as correctly processed.
2619+
/// Returns the number of messages negatively acknowledged.
2620+
/// </summary>
2621+
/// <param name="key">The key of the stream.</param>
2622+
/// <param name="groupName">The name of the consumer group that received the message.</param>
2623+
/// <param name="consumerName">The name of the consumer releasing the message.</param>
2624+
/// <param name="mode">The negative acknowledge mode to use.</param>
2625+
/// <param name="messageId">The ID of the message to negatively acknowledge.</param>
2626+
/// <param name="flags">The flags to use for this operation.</param>
2627+
/// <returns>The number of messages negatively acknowledged.</returns>
2628+
/// <remarks><seealso href="https://redis.io/topics/streams-intro"/></remarks>
2629+
[Experimental(Experiments.Server_8_8, UrlFormat = Experiments.UrlFormat)]
2630+
long StreamNegativeAcknowledge(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue messageId, CommandFlags flags = CommandFlags.None);
2631+
2632+
/// <summary>
2633+
/// Allow the consumer to release pending messages back to the group without marking them as correctly processed.
2634+
/// Returns the number of messages negatively acknowledged.
2635+
/// </summary>
2636+
/// <param name="key">The key of the stream.</param>
2637+
/// <param name="groupName">The name of the consumer group that received the messages.</param>
2638+
/// <param name="consumerName">The name of the consumer releasing the messages.</param>
2639+
/// <param name="mode">The negative acknowledge mode to use.</param>
2640+
/// <param name="messageIds">The IDs of the messages to negatively acknowledge.</param>
2641+
/// <param name="flags">The flags to use for this operation.</param>
2642+
/// <returns>The number of messages negatively acknowledged.</returns>
2643+
/// <remarks><seealso href="https://redis.io/topics/streams-intro"/></remarks>
2644+
[Experimental(Experiments.Server_8_8, UrlFormat = Experiments.UrlFormat)]
2645+
long StreamNegativeAcknowledge(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None);
2646+
26172647
/// <summary>
26182648
/// Adds an entry using the specified values to the given stream key.
26192649
/// If key does not exist, a new key holding a stream is created.

src/StackExchange.Redis/Interfaces/IDatabaseAsync.cs

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -644,6 +644,14 @@ IAsyncEnumerable<SortedSetEntry> SortedSetScanAsync(
644644
Task<StreamTrimResult[]> StreamAcknowledgeAndDeleteAsync(RedisKey key, RedisValue groupName, StreamTrimMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None);
645645
#pragma warning restore RS0026
646646

647+
/// <inheritdoc cref="IDatabase.StreamNegativeAcknowledge(RedisKey, RedisValue, RedisValue, StreamNackMode, RedisValue, CommandFlags)"/>
648+
[Experimental(Experiments.Server_8_8, UrlFormat = Experiments.UrlFormat)]
649+
Task<long> StreamNegativeAcknowledgeAsync(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue messageId, CommandFlags flags = CommandFlags.None);
650+
651+
/// <inheritdoc cref="IDatabase.StreamNegativeAcknowledge(RedisKey, RedisValue, RedisValue, StreamNackMode, RedisValue[], CommandFlags)"/>
652+
[Experimental(Experiments.Server_8_8, UrlFormat = Experiments.UrlFormat)]
653+
Task<long> StreamNegativeAcknowledgeAsync(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None);
654+
647655
/// <inheritdoc cref="IDatabase.StreamAdd(RedisKey, RedisValue, RedisValue, RedisValue?, int?, bool, CommandFlags)"/>
648656
Task<RedisValue> StreamAddAsync(RedisKey key, RedisValue streamField, RedisValue streamValue, RedisValue? messageId, int? maxLength, bool useApproximateMaxLength, CommandFlags flags);
649657

src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixed.cs

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -609,6 +609,12 @@ public Task<StreamTrimResult> StreamAcknowledgeAndDeleteAsync(RedisKey key, Redi
609609
public Task<StreamTrimResult[]> StreamAcknowledgeAndDeleteAsync(RedisKey key, RedisValue groupName, StreamTrimMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None) =>
610610
Inner.StreamAcknowledgeAndDeleteAsync(ToInner(key), groupName, mode, messageIds, flags);
611611

612+
public Task<long> StreamNegativeAcknowledgeAsync(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue messageId, CommandFlags flags = CommandFlags.None) =>
613+
Inner.StreamNegativeAcknowledgeAsync(ToInner(key), groupName, consumerName, mode, messageId, flags);
614+
615+
public Task<long> StreamNegativeAcknowledgeAsync(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None) =>
616+
Inner.StreamNegativeAcknowledgeAsync(ToInner(key), groupName, consumerName, mode, messageIds, flags);
617+
612618
public Task<RedisValue> StreamAddAsync(RedisKey key, RedisValue streamField, RedisValue streamValue, RedisValue? messageId, int? maxLength, bool useApproximateMaxLength, CommandFlags flags) =>
613619
Inner.StreamAddAsync(ToInner(key), streamField, streamValue, messageId, maxLength, useApproximateMaxLength, flags);
614620

src/StackExchange.Redis/KeyspaceIsolation/KeyPrefixedDatabase.cs

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -591,6 +591,12 @@ public StreamTrimResult StreamAcknowledgeAndDelete(RedisKey key, RedisValue grou
591591
public StreamTrimResult[] StreamAcknowledgeAndDelete(RedisKey key, RedisValue groupName, StreamTrimMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None) =>
592592
Inner.StreamAcknowledgeAndDelete(ToInner(key), groupName, mode, messageIds, flags);
593593

594+
public long StreamNegativeAcknowledge(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue messageId, CommandFlags flags = CommandFlags.None) =>
595+
Inner.StreamNegativeAcknowledge(ToInner(key), groupName, consumerName, mode, messageId, flags);
596+
597+
public long StreamNegativeAcknowledge(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None) =>
598+
Inner.StreamNegativeAcknowledge(ToInner(key), groupName, consumerName, mode, messageIds, flags);
599+
594600
public RedisValue StreamAdd(RedisKey key, RedisValue streamField, RedisValue streamValue, RedisValue? messageId, int? maxLength, bool useApproximateMaxLength, CommandFlags flags) =>
595601
Inner.StreamAdd(ToInner(key), streamField, streamValue, messageId, maxLength, useApproximateMaxLength, flags);
596602

src/StackExchange.Redis/PublicAPI/PublicAPI.Shipped.txt

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2259,9 +2259,13 @@ static StackExchange.Redis.RedisChannel.KeySpaceSingleKey(in StackExchange.Redis
22592259
[SER003]StackExchange.Redis.IDatabase.StreamAdd(StackExchange.Redis.RedisKey key, StackExchange.Redis.NameValueEntry[]! streamPairs, StackExchange.Redis.StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StackExchange.Redis.StreamTrimMode trimMode = StackExchange.Redis.StreamTrimMode.KeepReferences, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> StackExchange.Redis.RedisValue
22602260
[SER003]StackExchange.Redis.IDatabase.StreamAdd(StackExchange.Redis.RedisKey key, StackExchange.Redis.RedisValue streamField, StackExchange.Redis.RedisValue streamValue, StackExchange.Redis.StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StackExchange.Redis.StreamTrimMode trimMode = StackExchange.Redis.StreamTrimMode.KeepReferences, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> StackExchange.Redis.RedisValue
22612261
[SER003]StackExchange.Redis.IDatabase.StreamConfigure(StackExchange.Redis.RedisKey key, StackExchange.Redis.StreamConfiguration! configuration, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> void
2262+
[SER006]StackExchange.Redis.IDatabase.StreamNegativeAcknowledge(StackExchange.Redis.RedisKey key, StackExchange.Redis.RedisValue groupName, StackExchange.Redis.RedisValue consumerName, StackExchange.Redis.StreamNackMode mode, StackExchange.Redis.RedisValue messageId, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> long
2263+
[SER006]StackExchange.Redis.IDatabase.StreamNegativeAcknowledge(StackExchange.Redis.RedisKey key, StackExchange.Redis.RedisValue groupName, StackExchange.Redis.RedisValue consumerName, StackExchange.Redis.StreamNackMode mode, StackExchange.Redis.RedisValue[]! messageIds, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> long
22622264
[SER003]StackExchange.Redis.IDatabaseAsync.StreamAddAsync(StackExchange.Redis.RedisKey key, StackExchange.Redis.NameValueEntry[]! streamPairs, StackExchange.Redis.StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StackExchange.Redis.StreamTrimMode trimMode = StackExchange.Redis.StreamTrimMode.KeepReferences, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> System.Threading.Tasks.Task<StackExchange.Redis.RedisValue>!
22632265
[SER003]StackExchange.Redis.IDatabaseAsync.StreamAddAsync(StackExchange.Redis.RedisKey key, StackExchange.Redis.RedisValue streamField, StackExchange.Redis.RedisValue streamValue, StackExchange.Redis.StreamIdempotentId idempotentId, long? maxLength = null, bool useApproximateMaxLength = false, long? limit = null, StackExchange.Redis.StreamTrimMode trimMode = StackExchange.Redis.StreamTrimMode.KeepReferences, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> System.Threading.Tasks.Task<StackExchange.Redis.RedisValue>!
22642266
[SER003]StackExchange.Redis.IDatabaseAsync.StreamConfigureAsync(StackExchange.Redis.RedisKey key, StackExchange.Redis.StreamConfiguration! configuration, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> System.Threading.Tasks.Task!
2267+
[SER006]StackExchange.Redis.IDatabaseAsync.StreamNegativeAcknowledgeAsync(StackExchange.Redis.RedisKey key, StackExchange.Redis.RedisValue groupName, StackExchange.Redis.RedisValue consumerName, StackExchange.Redis.StreamNackMode mode, StackExchange.Redis.RedisValue messageId, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> System.Threading.Tasks.Task<long>!
2268+
[SER006]StackExchange.Redis.IDatabaseAsync.StreamNegativeAcknowledgeAsync(StackExchange.Redis.RedisKey key, StackExchange.Redis.RedisValue groupName, StackExchange.Redis.RedisValue consumerName, StackExchange.Redis.StreamNackMode mode, StackExchange.Redis.RedisValue[]! messageIds, StackExchange.Redis.CommandFlags flags = StackExchange.Redis.CommandFlags.None) -> System.Threading.Tasks.Task<long>!
22652269
[SER003]StackExchange.Redis.StreamConfiguration
22662270
[SER003]StackExchange.Redis.StreamConfiguration.IdmpDuration.get -> long?
22672271
[SER003]StackExchange.Redis.StreamConfiguration.IdmpDuration.set -> void
@@ -2280,6 +2284,10 @@ static StackExchange.Redis.RedisChannel.KeySpaceSingleKey(in StackExchange.Redis
22802284
[SER003]StackExchange.Redis.StreamInfo.IidsDuplicates.get -> long
22812285
[SER003]StackExchange.Redis.StreamInfo.IidsTracked.get -> long
22822286
[SER003]StackExchange.Redis.StreamInfo.PidsTracked.get -> long
2287+
[SER006]StackExchange.Redis.StreamNackMode
2288+
[SER006]StackExchange.Redis.StreamNackMode.Fail = 1 -> StackExchange.Redis.StreamNackMode
2289+
[SER006]StackExchange.Redis.StreamNackMode.Fatal = 2 -> StackExchange.Redis.StreamNackMode
2290+
[SER006]StackExchange.Redis.StreamNackMode.Silent = 0 -> StackExchange.Redis.StreamNackMode
22832291
StackExchange.Redis.StreamInfo.EntriesAdded.get -> long
22842292
StackExchange.Redis.StreamInfo.MaxDeletedEntryId.get -> StackExchange.Redis.RedisValue
22852293
StackExchange.Redis.StreamInfo.RecordedFirstEntryId.get -> StackExchange.Redis.RedisValue
Lines changed: 130 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,130 @@
1+
using System;
2+
using System.Threading.Tasks;
3+
4+
namespace StackExchange.Redis;
5+
6+
internal partial class RedisDatabase
7+
{
8+
public long StreamNegativeAcknowledge(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue messageId, CommandFlags flags = CommandFlags.None)
9+
=> ExecuteSync(GetStreamNegativeAcknowledgeMessage(key, groupName, consumerName, mode, messageId, flags), ResultProcessor.Int64);
10+
11+
public Task<long> StreamNegativeAcknowledgeAsync(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue messageId, CommandFlags flags = CommandFlags.None)
12+
=> ExecuteAsync(GetStreamNegativeAcknowledgeMessage(key, groupName, consumerName, mode, messageId, flags), ResultProcessor.Int64);
13+
14+
public long StreamNegativeAcknowledge(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None)
15+
=> ExecuteSync(GetStreamNegativeAcknowledgeMessage(key, groupName, consumerName, mode, messageIds, flags), ResultProcessor.Int64);
16+
17+
public Task<long> StreamNegativeAcknowledgeAsync(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue[] messageIds, CommandFlags flags = CommandFlags.None)
18+
=> ExecuteAsync(GetStreamNegativeAcknowledgeMessage(key, groupName, consumerName, mode, messageIds, flags), ResultProcessor.Int64);
19+
20+
private Message GetStreamNegativeAcknowledgeMessage(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue messageId, CommandFlags flags)
21+
=> new StreamNackMessageSingle(Database, flags, key, groupName, consumerName, mode, messageId);
22+
23+
private Message GetStreamNegativeAcknowledgeMessage(RedisKey key, RedisValue groupName, RedisValue consumerName, StreamNackMode mode, RedisValue[] messageIds, CommandFlags flags)
24+
=> messageIds is { Length: 1 }
25+
? new StreamNackMessageSingle(Database, flags, key, groupName, consumerName, mode, messageIds[0])
26+
: new StreamNackMessageMulti(Database, flags, key, groupName, consumerName, mode, messageIds);
27+
28+
internal abstract class StreamNackMessageBase : Message.CommandKeyBase
29+
{
30+
private readonly RedisValue groupName;
31+
private readonly RedisValue consumerName;
32+
private readonly StreamNackMode mode;
33+
34+
protected StreamNackMessageBase(int db, CommandFlags flags, in RedisKey key, in RedisValue groupName, in RedisValue consumerName, StreamNackMode mode)
35+
: base(db, flags, RedisCommand.XNACK, key)
36+
{
37+
groupName.AssertNotNull();
38+
consumerName.AssertNotNull();
39+
40+
this.groupName = groupName;
41+
this.consumerName = consumerName;
42+
this.mode = mode;
43+
}
44+
45+
protected abstract int Count { get; }
46+
47+
protected abstract void WriteIds(PhysicalConnection physical);
48+
49+
protected override void WriteImpl(PhysicalConnection physical)
50+
{
51+
physical.WriteHeader(Command, ArgCount);
52+
physical.Write(Key);
53+
physical.WriteBulkString(groupName);
54+
physical.WriteBulkString(consumerName);
55+
WriteMode(physical);
56+
physical.WriteBulkString(StreamConstants.Ids);
57+
physical.WriteBulkString(Count);
58+
WriteIds(physical);
59+
}
60+
61+
private void WriteMode(PhysicalConnection physical)
62+
{
63+
switch (mode)
64+
{
65+
case StreamNackMode.Silent:
66+
physical.WriteBulkString("SILENT"u8);
67+
break;
68+
case StreamNackMode.Fail:
69+
physical.WriteBulkString("FAIL"u8);
70+
break;
71+
case StreamNackMode.Fatal:
72+
physical.WriteBulkString("FATAL"u8);
73+
break;
74+
default:
75+
throw new ArgumentOutOfRangeException(nameof(mode));
76+
}
77+
}
78+
79+
public override int ArgCount => 6 + Count;
80+
}
81+
82+
internal sealed class StreamNackMessageSingle : StreamNackMessageBase
83+
{
84+
private readonly RedisValue messageId;
85+
86+
public StreamNackMessageSingle(int db, CommandFlags flags, in RedisKey key, in RedisValue groupName, in RedisValue consumerName, StreamNackMode mode, in RedisValue messageId)
87+
: base(db, flags, key, groupName, consumerName, mode)
88+
{
89+
messageId.AssertNotNull();
90+
this.messageId = messageId;
91+
}
92+
93+
protected override int Count => 1;
94+
95+
protected override void WriteIds(PhysicalConnection physical) => physical.WriteBulkString(messageId);
96+
}
97+
98+
internal sealed class StreamNackMessageMulti : StreamNackMessageBase
99+
{
100+
private readonly RedisValue[] messageIds;
101+
102+
public StreamNackMessageMulti(int db, CommandFlags flags, in RedisKey key, in RedisValue groupName, in RedisValue consumerName, StreamNackMode mode, RedisValue[] messageIds)
103+
: base(db, flags, key, groupName, consumerName, mode)
104+
{
105+
#if NET
106+
ArgumentNullException.ThrowIfNull(messageIds);
107+
#else
108+
if (messageIds == null) throw new ArgumentNullException(nameof(messageIds));
109+
#endif
110+
if (messageIds.Length == 0) throw new ArgumentOutOfRangeException(nameof(messageIds), "messageIds must contain at least one item.");
111+
112+
for (int i = 0; i < messageIds.Length; i++)
113+
{
114+
messageIds[i].AssertNotNull();
115+
}
116+
117+
this.messageIds = messageIds;
118+
}
119+
120+
protected override int Count => messageIds.Length;
121+
122+
protected override void WriteIds(PhysicalConnection physical)
123+
{
124+
for (int i = 0; i < messageIds.Length; i++)
125+
{
126+
physical.WriteBulkString(messageIds[i]);
127+
}
128+
}
129+
}
130+
}

0 commit comments

Comments
 (0)