This is an automated email from the ASF dual-hosted git repository.

Jackie-Jiang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git


The following commit(s) were added to refs/heads/master by this push:
     new 931c67a0f25 Introduce reason codes for segment completion requests 
(#19251)
931c67a0f25 is described below

commit 931c67a0f2574568d5d82e1536997eb3f753291e
Author: Arunkumar Ganesan <[email protected]>
AuthorDate: Sat Oct 3 17:02:55 2026 -0400

    Introduce reason codes for segment completion requests (#19251)
---
 .../protocols/SegmentCompletionProtocol.java       | 96 +++++++++++++++++++++-
 .../protocols/SegmentCompletionProtocolTest.java   | 89 +++++++++++++++++++-
 .../resources/LLCSegmentCompletionHandlers.java    | 23 ++++--
 .../LLCSegmentCompletionHandlersTest.java          | 92 +++++++++++++++++++++
 .../helix/core/realtime/SegmentCompletionTest.java | 38 +++++++++
 .../realtime/RealtimeSegmentDataManager.java       | 20 ++---
 .../realtime/RealtimeSegmentDataManagerTest.java   | 10 ++-
 7 files changed, 345 insertions(+), 23 deletions(-)

diff --git 
a/pinot-common/src/main/java/org/apache/pinot/common/protocols/SegmentCompletionProtocol.java
 
b/pinot-common/src/main/java/org/apache/pinot/common/protocols/SegmentCompletionProtocol.java
index 797681a64a0..2788ead958d 100644
--- 
a/pinot-common/src/main/java/org/apache/pinot/common/protocols/SegmentCompletionProtocol.java
+++ 
b/pinot-common/src/main/java/org/apache/pinot/common/protocols/SegmentCompletionProtocol.java
@@ -26,6 +26,7 @@ import java.io.IOException;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.concurrent.TimeUnit;
+import javax.annotation.Nullable;
 import org.apache.pinot.common.utils.URIUtils;
 import org.apache.pinot.spi.utils.JsonUtils;
 
@@ -129,6 +130,7 @@ public class SegmentCompletionProtocol {
   public static final String PARAM_MEMORY_USED_BYTES = "memoryUsedBytes";
   public static final String PARAM_SEGMENT_SIZE_BYTES = "segmentSizeBytes";
   public static final String PARAM_REASON = "reason";
+  public static final String PARAM_REASON_CODE = "reasonCode";
   // Sent by servers to request additional time to build
   public static final String PARAM_EXTRA_TIME_SEC = "extraTimeSec";
   // Sent by servers to indicate the number of rows read so far
@@ -150,6 +152,62 @@ public class SegmentCompletionProtocol {
   // (like size reaching close to its limit or number of col values for a col 
is about to overflow int max)
   public static final String REASON_INDEX_CAPACITY_THRESHOLD_BREACHED = 
"indexCapacityThresholdBreached";
 
+  /// Stable wire codes for server stop reasons.
+  ///
+  /// Numeric IDs are part of the segment-completion protocol and must remain 
unique, stable, and non-reusable.
+  /// Controllers ignore unknown numeric `reasonCode` values and fall back to 
the retained legacy `reason` value.
+  /// When both a known `reasonCode` and a legacy `reason` are present, the 
numeric code takes precedence and
+  /// normalizes the legacy reason string carried through the request params.
+  public enum ReasonCode {
+    ROW_LIMIT(100, REASON_ROW_LIMIT),
+    TIME_LIMIT(110, REASON_TIME_LIMIT),
+    END_OF_PARTITION_GROUP(120, REASON_END_OF_PARTITION_GROUP),
+    FORCE_COMMIT_MESSAGE_RECEIVED(130, REASON_FORCE_COMMIT_MESSAGE_RECEIVED),
+    INDEX_CAPACITY_THRESHOLD_BREACHED(140, 
REASON_INDEX_CAPACITY_THRESHOLD_BREACHED);
+
+    private static final Map<Integer, ReasonCode> BY_ID = new HashMap<>();
+
+    static {
+      for (ReasonCode value : values()) {
+        BY_ID.put(value.getId(), value);
+      }
+    }
+
+    private final int _id;
+    private final String _reason;
+
+    ReasonCode(int id, String reason) {
+      _id = id;
+      _reason = reason;
+    }
+
+    public int getId() {
+      return _id;
+    }
+
+    public String getReason() {
+      return _reason;
+    }
+
+    @Nullable
+    public static ReasonCode fromCode(int reasonCode) {
+      return BY_ID.get(reasonCode);
+    }
+
+    @Nullable
+    public static ReasonCode fromReason(@Nullable String reason) {
+      if (reason == null) {
+        return null;
+      }
+      for (ReasonCode value : values()) {
+        if (value.getReason().equals(reason)) {
+          return value;
+        }
+      }
+      return null;
+    }
+  }
+
   // Canned responses
   public static final Response RESP_NOT_LEADER =
       new Response(new 
Response.Params().withStatus(ControllerResponseStatus.NOT_LEADER));
@@ -201,6 +259,9 @@ public class SegmentCompletionProtocol {
       if (_params.getReason() != null) {
         params.put(PARAM_REASON, _params.getReason());
       }
+      if (_params.getReasonCode() != null) {
+        params.put(PARAM_REASON_CODE, 
String.valueOf(_params.getReasonCode().getId()));
+      }
       if (_params.getBuildTimeMillis() > 0) {
         params.put(PARAM_BUILD_TIME_MILLIS, 
String.valueOf(_params.getBuildTimeMillis()));
       }
@@ -231,7 +292,10 @@ public class SegmentCompletionProtocol {
     public static class Params {
       private String _segmentName;
       private String _instanceId;
+      @Nullable
       private String _reason;
+      @Nullable
+      private ReasonCode _reasonCode;
       private int _numRows;
       private long _buildTimeMillis;
       private long _waitTimeMillis;
@@ -253,6 +317,7 @@ public class SegmentCompletionProtocol {
         _segmentSizeBytes = SEGMENT_SIZE_BYTES_DEFAULT;
         _streamPartitionMsgOffset = null;
         _reason = null;
+        _reasonCode = null;
       }
 
       public Params(Params params) {
@@ -267,6 +332,7 @@ public class SegmentCompletionProtocol {
         _segmentSizeBytes = params.getSegmentSizeBytes();
         _streamPartitionMsgOffset = params.getStreamPartitionMsgOffset();
         _reason = params.getReason();
+        _reasonCode = params.getReasonCode();
       }
 
       public Params withSegmentName(String segmentName) {
@@ -279,8 +345,29 @@ public class SegmentCompletionProtocol {
         return this;
       }
 
-      public Params withReason(String reason) {
+      public Params withReason(@Nullable String reason) {
         _reason = reason;
+        _reasonCode = ReasonCode.fromReason(reason);
+        return this;
+      }
+
+      /// Resolves the `reasonCode` query parameter. Missing or unknown values 
are ignored so the retained legacy
+      /// `reason` can continue to be used as fallback.
+      public Params withReasonCodeParam(@Nullable Integer reasonCode) {
+        ReasonCode parsedReasonCode = reasonCode != null ? 
ReasonCode.fromCode(reasonCode) : null;
+        if (parsedReasonCode != null) {
+          return withReasonCode(parsedReasonCode);
+        }
+        return this;
+      }
+
+      /// Sets the typed stop reason. A known code also sets the legacy reason 
string so older controllers can ignore
+      /// `reasonCode` and still use `reason`. Passing `null` clears only the 
typed code.
+      public Params withReasonCode(@Nullable ReasonCode reasonCode) {
+        _reasonCode = reasonCode;
+        if (reasonCode != null) {
+          _reason = reasonCode.getReason();
+        }
         return this;
       }
 
@@ -328,10 +415,16 @@ public class SegmentCompletionProtocol {
         return _segmentName;
       }
 
+      @Nullable
       public String getReason() {
         return _reason;
       }
 
+      @Nullable
+      public ReasonCode getReasonCode() {
+        return _reasonCode;
+      }
+
       public String getInstanceId() {
         return _instanceId;
       }
@@ -372,6 +465,7 @@ public class SegmentCompletionProtocol {
         return "Segment name: " + _segmentName
             + ",Instance Id: " + _instanceId
             + ",Reason: " + _reason
+            + ",ReasonCode: " + _reasonCode
             + ",NumRows: " + _numRows
             + ",BuildTimeMillis: " + _buildTimeMillis
             + ",WaitTimeMillis: " + _waitTimeMillis
diff --git 
a/pinot-common/src/test/java/org/apache/pinot/common/protocols/SegmentCompletionProtocolTest.java
 
b/pinot-common/src/test/java/org/apache/pinot/common/protocols/SegmentCompletionProtocolTest.java
index f5854cdb060..15483b220c8 100644
--- 
a/pinot-common/src/test/java/org/apache/pinot/common/protocols/SegmentCompletionProtocolTest.java
+++ 
b/pinot-common/src/test/java/org/apache/pinot/common/protocols/SegmentCompletionProtocolTest.java
@@ -20,7 +20,9 @@ package org.apache.pinot.common.protocols;
 
 import java.net.URI;
 import java.util.Arrays;
+import java.util.HashSet;
 import java.util.Map;
+import java.util.Set;
 import java.util.stream.Collectors;
 import org.apache.pinot.common.utils.URIUtils;
 import org.testng.Assert;
@@ -47,6 +49,7 @@ public class SegmentCompletionProtocolTest {
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_SEGMENT_NAME),
 "UNKNOWN_SEGMENT");
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_INSTANCE_ID), 
"UNKNOWN_INSTANCE");
     Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON));
+    
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE));
     
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS));
     
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS));
     
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_EXTRA_TIME_SEC));
@@ -69,6 +72,7 @@ public class SegmentCompletionProtocolTest {
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_SEGMENT_NAME),
 "foo__0__0__12345Z");
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_INSTANCE_ID), 
"Server_localhost_8099");
     Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON));
+    
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE));
     
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS));
     
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS));
     
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_EXTRA_TIME_SEC));
@@ -79,7 +83,8 @@ public class SegmentCompletionProtocolTest {
     
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_STREAM_PARTITION_MSG_OFFSET));
 
     params = new 
SegmentCompletionProtocol.Request.Params().withSegmentName("foo__0__0__12345Z")
-        
.withInstanceId("Server_localhost_8099").withReason("ROW_LIMIT").withBuildTimeMillis(1000)
+        
.withInstanceId("Server_localhost_8099").withReason(SegmentCompletionProtocol.REASON_ROW_LIMIT)
+        .withBuildTimeMillis(1000)
         
.withWaitTimeMillis(2000).withExtraTimeSec(3000).withMemoryUsedBytes(4000).withSegmentSizeBytes(5000)
         
.withNumRows(6000).withSegmentLocation("/tmp/segment").withStreamPartitionMsgOffset("7000");
     SegmentCompletionProtocol.SegmentCommitStartRequest 
commitStartRequestWithAllParams =
@@ -89,7 +94,9 @@ public class SegmentCompletionProtocolTest {
         Arrays.stream(uri.getQuery().split("&")).collect(Collectors.toMap(e -> 
e.split("=")[0], e -> e.split("=")[1]));
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_SEGMENT_NAME),
 "foo__0__0__12345Z");
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_INSTANCE_ID), 
"Server_localhost_8099");
-    Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON), 
"ROW_LIMIT");
+    Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON),
+        SegmentCompletionProtocol.REASON_ROW_LIMIT);
+    
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE), 
"100");
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS),
 "1000");
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS),
 "2000");
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_EXTRA_TIME_SEC),
 "3000");
@@ -122,6 +129,7 @@ public class SegmentCompletionProtocolTest {
         URIUtils.encode("Server_localhost_8099"));
     Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON),
         "%7B%22type%22%3A%22ROW_LIMIT%22%2C%20%22value%22%3A1000%7D");
+    
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE));
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS),
 "1000");
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS),
 "2000");
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_EXTRA_TIME_SEC),
 "3000");
@@ -133,4 +141,81 @@ public class SegmentCompletionProtocolTest {
     
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_STREAM_PARTITION_MSG_OFFSET),
         
URIUtils.encode("{\"shardId-000000000001\":\"49615238429973311938200772279310862572716999467690098706\"}"));
   }
+
+  @Test
+  public void testReasonCodeCompatibilityAndPrecedence() {
+    SegmentCompletionProtocol.Request.Params params =
+        new 
SegmentCompletionProtocol.Request.Params().withReasonCode(SegmentCompletionProtocol.ReasonCode.TIME_LIMIT);
+    Assert.assertEquals(params.getReason(), 
SegmentCompletionProtocol.REASON_TIME_LIMIT);
+    Assert.assertEquals(params.getReasonCode(), 
SegmentCompletionProtocol.ReasonCode.TIME_LIMIT);
+
+    params = new 
SegmentCompletionProtocol.Request.Params().withReason("customReason")
+        .withReasonCode(SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+    Assert.assertEquals(params.getReason(), 
SegmentCompletionProtocol.REASON_ROW_LIMIT);
+    Assert.assertEquals(params.getReasonCode(), 
SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+
+    params = new SegmentCompletionProtocol.Request.Params()
+        .withReasonCodeParam(110);
+    Assert.assertEquals(params.getReason(), 
SegmentCompletionProtocol.REASON_TIME_LIMIT);
+    Assert.assertEquals(params.getReasonCode(), 
SegmentCompletionProtocol.ReasonCode.TIME_LIMIT);
+
+    params = new 
SegmentCompletionProtocol.Request.Params().withReason(SegmentCompletionProtocol.REASON_ROW_LIMIT)
+        .withReasonCodeParam(999);
+    Assert.assertEquals(params.getReason(), 
SegmentCompletionProtocol.REASON_ROW_LIMIT);
+    Assert.assertEquals(params.getReasonCode(), 
SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+
+    params = new 
SegmentCompletionProtocol.Request.Params().withReasonCodeParam(999);
+    Assert.assertNull(params.getReason());
+    Assert.assertNull(params.getReasonCode());
+  }
+
+  @Test
+  public void testReasonCodeValuesAreStableAndUnique() {
+    Object[][] expectedValues = {
+        {SegmentCompletionProtocol.ReasonCode.ROW_LIMIT, 100, 
SegmentCompletionProtocol.REASON_ROW_LIMIT},
+        {SegmentCompletionProtocol.ReasonCode.TIME_LIMIT, 110, 
SegmentCompletionProtocol.REASON_TIME_LIMIT},
+        {SegmentCompletionProtocol.ReasonCode.END_OF_PARTITION_GROUP, 120,
+            SegmentCompletionProtocol.REASON_END_OF_PARTITION_GROUP},
+        {SegmentCompletionProtocol.ReasonCode.FORCE_COMMIT_MESSAGE_RECEIVED, 
130,
+            SegmentCompletionProtocol.REASON_FORCE_COMMIT_MESSAGE_RECEIVED},
+        
{SegmentCompletionProtocol.ReasonCode.INDEX_CAPACITY_THRESHOLD_BREACHED, 140,
+            SegmentCompletionProtocol.REASON_INDEX_CAPACITY_THRESHOLD_BREACHED}
+    };
+
+    Set<Integer> ids = new HashSet<>();
+    for (Object[] expectedValue : expectedValues) {
+      SegmentCompletionProtocol.ReasonCode reasonCode =
+          (SegmentCompletionProtocol.ReasonCode) expectedValue[0];
+      int id = (int) expectedValue[1];
+      String reason = (String) expectedValue[2];
+
+      Assert.assertTrue(ids.add(id), "Reason code IDs must be unique");
+      Assert.assertEquals(reasonCode.getId(), id);
+      Assert.assertEquals(reasonCode.getReason(), reason);
+      Assert.assertEquals(SegmentCompletionProtocol.ReasonCode.fromCode(id), 
reasonCode);
+      
Assert.assertEquals(SegmentCompletionProtocol.ReasonCode.fromReason(reason), 
reasonCode);
+    }
+
+    Assert.assertEquals(ids.size(), 
SegmentCompletionProtocol.ReasonCode.values().length);
+  }
+
+  @Test
+  public void testReasonCodeRequestRetainsLegacyReasonForOldController()
+      throws Exception {
+    SegmentCompletionProtocol.Request.Params params = new 
SegmentCompletionProtocol.Request.Params()
+        .withSegmentName("foo__0__0__12345Z")
+        .withInstanceId("Server_localhost_8099")
+        .withStreamPartitionMsgOffset("7000")
+        .withReasonCode(SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+    SegmentCompletionProtocol.SegmentConsumedRequest request =
+        new SegmentCompletionProtocol.SegmentConsumedRequest(params);
+
+    URI uri = new URI(request.getUrl("localhost:8080", "http"));
+    Map<String, String> paramsMap =
+        Arrays.stream(uri.getQuery().split("&")).collect(Collectors.toMap(e -> 
e.split("=")[0], e -> e.split("=")[1]));
+
+    Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON),
+        SegmentCompletionProtocol.REASON_ROW_LIMIT);
+    
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE), 
"100");
+  }
 }
diff --git 
a/pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlers.java
 
b/pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlers.java
index 4ebd5fa7de6..7f0cc8f7aec 100644
--- 
a/pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlers.java
+++ 
b/pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlers.java
@@ -122,6 +122,7 @@ public class LLCSegmentCompletionHandlers {
       @QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_NAME) String 
segmentName,
       @QueryParam(SegmentCompletionProtocol.PARAM_STREAM_PARTITION_MSG_OFFSET) 
String streamPartitionMsgOffset,
       @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason,
+      @QueryParam(SegmentCompletionProtocol.PARAM_REASON_CODE) Integer 
stopReasonCode,
       @QueryParam(SegmentCompletionProtocol.PARAM_MEMORY_USED_BYTES) long 
memoryUsedBytes,
       @QueryParam(SegmentCompletionProtocol.PARAM_ROW_COUNT) int numRows) {
     if (instanceId == null || segmentName == null || streamPartitionMsgOffset 
== null) {
@@ -135,6 +136,7 @@ public class LLCSegmentCompletionHandlers {
         .withSegmentName(segmentName)
         .withStreamPartitionMsgOffset(streamPartitionMsgOffset)
         .withReason(stopReason)
+        .withReasonCodeParam(stopReasonCode)
         .withMemoryUsedBytes(memoryUsedBytes)
         .withNumRows(numRows);
     LOGGER.info("Processing segmentConsumed: {}", requestParams);
@@ -152,7 +154,8 @@ public class LLCSegmentCompletionHandlers {
   public String 
segmentStoppedConsuming(@QueryParam(SegmentCompletionProtocol.PARAM_INSTANCE_ID)
 String instanceId,
       @QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_NAME) String 
segmentName,
       @QueryParam(SegmentCompletionProtocol.PARAM_STREAM_PARTITION_MSG_OFFSET) 
String streamPartitionMsgOffset,
-      @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason) {
+      @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason,
+      @QueryParam(SegmentCompletionProtocol.PARAM_REASON_CODE) Integer 
stopReasonCode) {
     if (instanceId == null || segmentName == null || streamPartitionMsgOffset 
== null) {
       LOGGER.error("Invalid call: segmentName={}, instanceId={}, 
streamPartitionMsgOffset={}", segmentName, instanceId,
           streamPartitionMsgOffset);
@@ -163,7 +166,8 @@ public class LLCSegmentCompletionHandlers {
         .withInstanceId(instanceId)
         .withSegmentName(segmentName)
         .withStreamPartitionMsgOffset(streamPartitionMsgOffset)
-        .withReason(stopReason);
+        .withReason(stopReason)
+        .withReasonCodeParam(stopReasonCode);
     LOGGER.info("Processing segmentStoppedConsuming: {}", requestParams);
 
     String response = 
_segmentCompletionManager.segmentStoppedConsuming(requestParams).toJsonString();
@@ -183,7 +187,9 @@ public class LLCSegmentCompletionHandlers {
       @QueryParam(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS) long 
buildTimeMillis,
       @QueryParam(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS) long 
waitTimeMillis,
       @QueryParam(SegmentCompletionProtocol.PARAM_ROW_COUNT) int numRows,
-      @QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_SIZE_BYTES) long 
segmentSizeBytes) {
+      @QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_SIZE_BYTES) long 
segmentSizeBytes,
+      @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason,
+      @QueryParam(SegmentCompletionProtocol.PARAM_REASON_CODE) Integer 
stopReasonCode) {
     if (instanceId == null || segmentName == null || streamPartitionMsgOffset 
== null) {
       LOGGER.error("Invalid call: segmentName={}, instanceId={}, 
streamPartitionMsgOffset={}", segmentName, instanceId,
           streamPartitionMsgOffset);
@@ -198,7 +204,9 @@ public class LLCSegmentCompletionHandlers {
         .withBuildTimeMillis(buildTimeMillis)
         .withWaitTimeMillis(waitTimeMillis)
         .withNumRows(numRows)
-        .withSegmentSizeBytes(segmentSizeBytes);
+        .withSegmentSizeBytes(segmentSizeBytes)
+        .withReason(stopReason)
+        .withReasonCodeParam(stopReasonCode);
     LOGGER.info("Processing segmentCommitStart: {}", requestParams);
 
     String response = 
_segmentCompletionManager.segmentCommitStart(requestParams).toJsonString();
@@ -273,7 +281,9 @@ public class LLCSegmentCompletionHandlers {
       @QueryParam(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS) long 
waitTimeMillis,
       @QueryParam(SegmentCompletionProtocol.PARAM_ROW_COUNT) int numRows,
       @QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_SIZE_BYTES) long 
segmentSizeBytes,
-      @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason, 
FormDataMultiPart metadataFiles) {
+      @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason,
+      @QueryParam(SegmentCompletionProtocol.PARAM_REASON_CODE) Integer 
stopReasonCode,
+      FormDataMultiPart metadataFiles) {
     if (instanceId == null || segmentName == null || segmentLocation == null 
|| metadataFiles == null
         || streamPartitionMsgOffset == null) {
       LOGGER.error("Invalid call: segmentName={}, instanceId={}, 
segmentLocation={}, streamPartitionMsgOffset={}",
@@ -292,7 +302,8 @@ public class LLCSegmentCompletionHandlers {
         .withWaitTimeMillis(waitTimeMillis)
         .withNumRows(numRows)
         .withMemoryUsedBytes(memoryUsedBytes)
-        .withReason(stopReason);
+        .withReason(stopReason)
+        .withReasonCodeParam(stopReasonCode);
     LOGGER.info("Processing segmentCommitEndWithMetadata: {}", requestParams);
 
     SegmentMetadataImpl segmentMetadata;
diff --git 
a/pinot-controller/src/test/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlersTest.java
 
b/pinot-controller/src/test/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlersTest.java
new file mode 100644
index 00000000000..83bf54cdda2
--- /dev/null
+++ 
b/pinot-controller/src/test/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlersTest.java
@@ -0,0 +1,92 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.controller.api.resources;
+
+import org.apache.pinot.common.protocols.SegmentCompletionProtocol;
+import 
org.apache.pinot.controller.helix.core.realtime.SegmentCompletionManager;
+import org.mockito.ArgumentCaptor;
+import org.testng.annotations.Test;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+import static org.testng.Assert.assertEquals;
+
+
+/// Tests reason-code binding at the segment-completion REST resource boundary.
+public class LLCSegmentCompletionHandlersTest {
+  private static final String INSTANCE_ID = "Server_localhost_8099";
+  private static final String SEGMENT_NAME = "foo__0__0__12345Z";
+  private static final String OFFSET = "7000";
+
+  @Test
+  public void testSegmentCommitStartBindsReasonAndReasonCode() {
+    SegmentCompletionManager segmentCompletionManager = 
mock(SegmentCompletionManager.class);
+    
when(segmentCompletionManager.segmentCommitStart(any())).thenReturn(SegmentCompletionProtocol.RESP_COMMIT_CONTINUE);
+    LLCSegmentCompletionHandlers handler = new LLCSegmentCompletionHandlers();
+    handler._segmentCompletionManager = segmentCompletionManager;
+
+    handler.segmentCommitStart(INSTANCE_ID, SEGMENT_NAME, OFFSET, 4000, 1000, 
2000, 6000, 5000, "legacyReason", 100);
+
+    SegmentCompletionProtocol.Request.Params params = 
captureSegmentCommitStartParams(segmentCompletionManager);
+    assertEquals(params.getReason(), 
SegmentCompletionProtocol.REASON_ROW_LIMIT);
+    assertEquals(params.getReasonCode(), 
SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+  }
+
+  @Test
+  public void testSegmentCommitStartFallsBackToLegacyReasonForUnknownCode() {
+    SegmentCompletionManager segmentCompletionManager = 
mock(SegmentCompletionManager.class);
+    
when(segmentCompletionManager.segmentCommitStart(any())).thenReturn(SegmentCompletionProtocol.RESP_COMMIT_CONTINUE);
+    LLCSegmentCompletionHandlers handler = new LLCSegmentCompletionHandlers();
+    handler._segmentCompletionManager = segmentCompletionManager;
+
+    handler.segmentCommitStart(INSTANCE_ID, SEGMENT_NAME, OFFSET, 4000, 1000, 
2000, 6000, 5000,
+        SegmentCompletionProtocol.REASON_TIME_LIMIT, 999);
+
+    SegmentCompletionProtocol.Request.Params params = 
captureSegmentCommitStartParams(segmentCompletionManager);
+    assertEquals(params.getReason(), 
SegmentCompletionProtocol.REASON_TIME_LIMIT);
+    assertEquals(params.getReasonCode(), 
SegmentCompletionProtocol.ReasonCode.TIME_LIMIT);
+  }
+
+  @Test
+  public void testSegmentConsumedPrefersKnownReasonCodeOverLegacyReason() {
+    SegmentCompletionManager segmentCompletionManager = 
mock(SegmentCompletionManager.class);
+    
when(segmentCompletionManager.segmentConsumed(any())).thenReturn(SegmentCompletionProtocol.RESP_FAILED);
+    LLCSegmentCompletionHandlers handler = new LLCSegmentCompletionHandlers();
+    handler._segmentCompletionManager = segmentCompletionManager;
+
+    handler.segmentConsumed(INSTANCE_ID, SEGMENT_NAME, OFFSET, "legacyReason", 
100, 4000, 6000);
+
+    ArgumentCaptor<SegmentCompletionProtocol.Request.Params> paramsCaptor =
+        
ArgumentCaptor.forClass(SegmentCompletionProtocol.Request.Params.class);
+    verify(segmentCompletionManager).segmentConsumed(paramsCaptor.capture());
+    SegmentCompletionProtocol.Request.Params params = paramsCaptor.getValue();
+    assertEquals(params.getReason(), 
SegmentCompletionProtocol.REASON_ROW_LIMIT);
+    assertEquals(params.getReasonCode(), 
SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+  }
+
+  private static SegmentCompletionProtocol.Request.Params 
captureSegmentCommitStartParams(
+      SegmentCompletionManager segmentCompletionManager) {
+    ArgumentCaptor<SegmentCompletionProtocol.Request.Params> paramsCaptor =
+        
ArgumentCaptor.forClass(SegmentCompletionProtocol.Request.Params.class);
+    
verify(segmentCompletionManager).segmentCommitStart(paramsCaptor.capture());
+    return paramsCaptor.getValue();
+  }
+}
diff --git 
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/SegmentCompletionTest.java
 
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/SegmentCompletionTest.java
index 4c33fc249dd..95f80ee9a2a 100644
--- 
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/SegmentCompletionTest.java
+++ 
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/SegmentCompletionTest.java
@@ -589,6 +589,44 @@ public class SegmentCompletionTest {
     Assert.assertEquals(response.getStatus(), 
SegmentCompletionProtocol.ControllerResponseStatus.HOLD);
   }
 
+  @Test
+  public void testWinnerOnRowLimitReasonCodeOverridesLegacyReason() {
+    SegmentCompletionProtocol.Response response;
+    Request.Params params;
+    _segmentCompletionMgr._seconds = 10L;
+    params = new 
Request.Params().withInstanceId(S_1).withStreamPartitionMsgOffset(_s1Offset.toString())
+        .withSegmentName(_segmentNameStr)
+        .withReason("unknownReason")
+        .withReasonCode(SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+    response = _segmentCompletionMgr.segmentConsumed(params);
+    Assert.assertEquals(response.getStatus(), 
SegmentCompletionProtocol.ControllerResponseStatus.COMMIT);
+  }
+
+  @Test
+  public void testUnknownReasonCodeFallsBackToLegacyReason() {
+    SegmentCompletionProtocol.Response response;
+    Request.Params params;
+    _segmentCompletionMgr._seconds = 10L;
+    params = new 
Request.Params().withInstanceId(S_1).withStreamPartitionMsgOffset(_s1Offset.toString())
+        .withSegmentName(_segmentNameStr)
+        .withReason(SegmentCompletionProtocol.REASON_ROW_LIMIT)
+        .withReasonCodeParam(999);
+    response = _segmentCompletionMgr.segmentConsumed(params);
+    Assert.assertEquals(response.getStatus(), 
SegmentCompletionProtocol.ControllerResponseStatus.COMMIT);
+  }
+
+  @Test
+  public void testUnknownReasonCodeWithoutLegacyReason() {
+    SegmentCompletionProtocol.Response response;
+    Request.Params params;
+    _segmentCompletionMgr._seconds = 10L;
+    params = new 
Request.Params().withInstanceId(S_1).withStreamPartitionMsgOffset(_s1Offset.toString())
+        .withSegmentName(_segmentNameStr)
+        .withReasonCodeParam(999);
+    response = _segmentCompletionMgr.segmentConsumed(params);
+    Assert.assertEquals(response.getStatus(), 
SegmentCompletionProtocol.ControllerResponseStatus.HOLD);
+  }
+
   @Test
   public void testWinnerOnForceCommit() {
     SegmentCompletionProtocol.Response response;
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
index 1e239ce14da..f096ab6cc92 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
@@ -335,7 +335,7 @@ public class RealtimeSegmentDataManager extends 
SegmentDataManager {
   private long _consumeStartTime = -1;
   private long _lastLogTime = 0;
   private int _lastConsumedCount = 0;
-  private String _stopReason = null;
+  private SegmentCompletionProtocol.ReasonCode _stopReasonCode = null;
   private final Semaphore _segBuildSemaphore;
   private final boolean _isOffHeap;
   /// Whether null handling is enabled by default. This value is only used if
@@ -376,35 +376,35 @@ public class RealtimeSegmentDataManager extends 
SegmentDataManager {
                     _startTimeMs, now, _numRowsConsumed, _numRowsIndexed);
             _stopReasonPrinted = true;
           }
-          _stopReason = SegmentCompletionProtocol.REASON_TIME_LIMIT;
+          _stopReasonCode = SegmentCompletionProtocol.ReasonCode.TIME_LIMIT;
           return true;
         } else if (_numRowsIndexed >= _segmentMaxRowCount) {
           _segmentLogger.info("Stopping consumption due to row limit nRows={} 
numRowsIndexed={}, numRowsConsumed={}",
               _segmentMaxRowCount, _numRowsIndexed, _numRowsConsumed);
-          _stopReason = SegmentCompletionProtocol.REASON_ROW_LIMIT;
+          _stopReasonCode = SegmentCompletionProtocol.ReasonCode.ROW_LIMIT;
           return true;
         } else if (_endOfPartitionGroup) {
           _segmentLogger.info("Stopping consumption due to end of 
partitionGroup reached nRows={} numRowsIndexed={}, "
               + "numRowsConsumed={}", _segmentMaxRowCount, _numRowsIndexed, 
_numRowsConsumed);
-          _stopReason = 
SegmentCompletionProtocol.REASON_END_OF_PARTITION_GROUP;
+          _stopReasonCode = 
SegmentCompletionProtocol.ReasonCode.END_OF_PARTITION_GROUP;
           return true;
         } else if (_forceCommitMessageReceived) {
           _segmentLogger.info("Stopping consumption due to force commit - 
numRowsConsumed={} numRowsIndexed={}",
               _numRowsConsumed, _numRowsIndexed);
-          _stopReason = 
SegmentCompletionProtocol.REASON_FORCE_COMMIT_MESSAGE_RECEIVED;
+          _stopReasonCode = 
SegmentCompletionProtocol.ReasonCode.FORCE_COMMIT_MESSAGE_RECEIVED;
           return true;
         } else if (!canAddMore()) {
           _segmentLogger.info(
               "Stopping consumption as mutable index cannot consume more rows 
- numRowsConsumed={} "
                   + "numRowsIndexed={}",
               _numRowsConsumed, _numRowsIndexed);
-          _stopReason = 
SegmentCompletionProtocol.REASON_INDEX_CAPACITY_THRESHOLD_BREACHED;
+          _stopReasonCode = 
SegmentCompletionProtocol.ReasonCode.INDEX_CAPACITY_THRESHOLD_BREACHED;
           return true;
         }
         return false;
 
       case CATCHING_UP:
-        _stopReason = null;
+        _stopReasonCode = null;
         // We have posted segmentConsumed() at least once, and the controller 
is asking us to catch up to a certain
         // offset.
         // There is no time limit here, so just check to see that we are still 
within the offset we need to reach.
@@ -1071,7 +1071,7 @@ public class RealtimeSegmentDataManager extends 
SegmentDataManager {
   private boolean startSegmentCommit() {
     SegmentCompletionProtocol.Request.Params params = new 
SegmentCompletionProtocol.Request.Params();
     
params.withSegmentName(_segmentNameStr).withStreamPartitionMsgOffset(_currentOffset.toString())
-        
.withNumRows(_numRowsIndexed).withInstanceId(_instanceId).withReason(_stopReason);
+        
.withNumRows(_numRowsIndexed).withInstanceId(_instanceId).withReasonCode(_stopReasonCode);
     if (_isOffHeap) {
       params.withMemoryUsedBytes(_memoryManager.getTotalAllocatedBytes());
     }
@@ -1373,7 +1373,7 @@ public class RealtimeSegmentDataManager extends 
SegmentDataManager {
     SegmentCompletionProtocol.Request.Params params = new 
SegmentCompletionProtocol.Request.Params();
 
     
params.withSegmentName(_segmentNameStr).withStreamPartitionMsgOffset(_currentOffset.toString())
-        
.withNumRows(_numRowsIndexed).withInstanceId(_instanceId).withReason(_stopReason)
+        
.withNumRows(_numRowsIndexed).withInstanceId(_instanceId).withReasonCode(_stopReasonCode)
         .withBuildTimeMillis(_segmentBuildDescriptor.getBuildTimeMillis())
         .withSegmentSizeBytes(_segmentBuildDescriptor.getSegmentSizeBytes())
         .withWaitTimeMillis(_segmentBuildDescriptor.getWaitTimeMillis());
@@ -1610,7 +1610,7 @@ public class RealtimeSegmentDataManager extends 
SegmentDataManager {
     // Retry maybe once if leader is not found.
     SegmentCompletionProtocol.Request.Params params = new 
SegmentCompletionProtocol.Request.Params();
     
params.withStreamPartitionMsgOffset(_currentOffset.toString()).withSegmentName(_segmentNameStr)
-        
.withReason(_stopReason).withNumRows(_numRowsIndexed).withInstanceId(_instanceId);
+        
.withReasonCode(_stopReasonCode).withNumRows(_numRowsIndexed).withInstanceId(_instanceId);
     if (_isOffHeap) {
       params.withMemoryUsedBytes(_memoryManager.getTotalAllocatedBytes());
     }
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
 
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
index 7aec91dc0d2..c778794ffb8 100644
--- 
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
@@ -1491,7 +1491,7 @@ public class RealtimeSegmentDataManagerTest {
 
     public Field _state;
     public Field _shouldStop;
-    public Field _stopReason;
+    public Field _stopReasonCode;
     public Field _segmentBuildFailedWithDeterministicError;
     public boolean _failSegmentBuildAndReplace = false;
     // When set, buildSegmentAndReplace runs the real implementation 
(including the upsert CRC guard) instead of being
@@ -1545,8 +1545,8 @@ public class RealtimeSegmentDataManagerTest {
       _state.setAccessible(true);
       _shouldStop = 
RealtimeSegmentDataManager.class.getDeclaredField("_shouldStop");
       _shouldStop.setAccessible(true);
-      _stopReason = 
RealtimeSegmentDataManager.class.getDeclaredField("_stopReason");
-      _stopReason.setAccessible(true);
+      _stopReasonCode = 
RealtimeSegmentDataManager.class.getDeclaredField("_stopReasonCode");
+      _stopReasonCode.setAccessible(true);
       _segmentBuildFailedWithDeterministicError =
           
RealtimeSegmentDataManager.class.getDeclaredField("_segmentBuildFailedWithDeterministicError");
       _segmentBuildFailedWithDeterministicError.setAccessible(true);
@@ -1578,7 +1578,9 @@ public class RealtimeSegmentDataManagerTest {
 
     public String getStopReason() {
       try {
-        return (String) _stopReason.get(this);
+        SegmentCompletionProtocol.ReasonCode stopReasonCode =
+            (SegmentCompletionProtocol.ReasonCode) _stopReasonCode.get(this);
+        return stopReasonCode != null ? stopReasonCode.getReason() : null;
       } catch (Exception e) {
         Assert.fail();
       }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to