Skip to content

Commit 4e77472

Browse files
feat: introduce a generic map for per message nack options (#222)
Signed-off-by: Vaibhav Tiwari <vaibhav.tiwari33@gmail.com>
1 parent 50f3be8 commit 4e77472

3 files changed

Lines changed: 22 additions & 2 deletions

File tree

src/main/java/io/numaproj/numaflow/shared/NackOptions.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@
33
import lombok.Builder;
44
import lombok.Getter;
55

6+
import java.util.Map;
7+
68
/**
79
* NackOptions carries per-message redelivery options for a negative acknowledgement (nack).
810
* All fields are optional; a null value means unset.
@@ -16,6 +18,8 @@ public class NackOptions {
1618
private final Integer maxDeliveries;
1719
/** human-readable reason for the nack. */
1820
private final String reason;
21+
/** generic values passed as nack options */
22+
private final Map<String, String> nackMap;
1923

2024
/** Converts to the outgoing proto type, setting only the fields that are present. */
2125
public common.NackOptionsOuterClass.NackOptions toProto() {
@@ -30,6 +34,9 @@ public common.NackOptionsOuterClass.NackOptions toProto() {
3034
if (reason != null) {
3135
b.setReason(reason);
3236
}
37+
if (nackMap != null) {
38+
b.putAllNackMap(nackMap);
39+
}
3340
return b.build();
3441
}
3542

@@ -48,6 +55,9 @@ public static NackOptions fromProto(common.NackOptionsOuterClass.NackOptions p)
4855
if (p.hasReason()) {
4956
b.reason(p.getReason());
5057
}
58+
if (!p.getNackMapMap().isEmpty()) {
59+
b.nackMap(p.getNackMapMap());
60+
}
5161
return b.build();
5262
}
5363
}

src/main/proto/common/nack_options.proto

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,4 +6,5 @@ message NackOptions {
66
optional string reason = 1;
77
optional uint32 max_deliveries = 2;
88
optional uint64 delay = 3;
9+
map<string, string> nack_map = 4;
910
}

src/test/java/io/numaproj/numaflow/shared/NackOptionsTest.java

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,8 @@
22

33
import org.junit.Test;
44

5+
import java.util.HashMap;
6+
57
import static org.junit.Assert.assertEquals;
68
import static org.junit.Assert.assertFalse;
79
import static org.junit.Assert.assertNull;
@@ -11,14 +13,18 @@ public class NackOptionsTest {
1113

1214
@Test
1315
public void toProto_allFields() {
14-
NackOptions n = NackOptions.newBuilder().delay(500L).maxDeliveries(3).reason("retry").build();
16+
HashMap<String, String> nackMap = new HashMap<>();
17+
nackMap.put("key", "value");
18+
NackOptions n = NackOptions.newBuilder().delay(500L).maxDeliveries(3).reason("retry").nackMap(nackMap).build();
1519
common.NackOptionsOuterClass.NackOptions p = n.toProto();
1620
assertTrue(p.hasDelay());
1721
assertEquals(500L, p.getDelay());
1822
assertTrue(p.hasMaxDeliveries());
1923
assertEquals(3, p.getMaxDeliveries());
2024
assertTrue(p.hasReason());
2125
assertEquals("retry", p.getReason());
26+
assertFalse(p.getNackMapMap().isEmpty());
27+
assertEquals("value", p.getNackMapMap().get("key"));
2228
}
2329

2430
@Test
@@ -32,12 +38,15 @@ public void toProto_partialFields() {
3238

3339
@Test
3440
public void fromProto_roundTrip() {
41+
HashMap<String, String> nackMap = new HashMap<>();
42+
nackMap.put("key", "value");
3543
common.NackOptionsOuterClass.NackOptions p = common.NackOptionsOuterClass.NackOptions.newBuilder()
36-
.setDelay(500L).setMaxDeliveries(3).setReason("retry").build();
44+
.setDelay(500L).setMaxDeliveries(3).setReason("retry").putAllNackMap(nackMap).build();
3745
NackOptions n = NackOptions.fromProto(p);
3846
assertEquals(Long.valueOf(500L), n.getDelay());
3947
assertEquals(Integer.valueOf(3), n.getMaxDeliveries());
4048
assertEquals("retry", n.getReason());
49+
assertEquals(nackMap, n.getNackMap());
4150
}
4251

4352
@Test

0 commit comments

Comments
 (0)