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

errose28 pushed a commit to branch HDDS-14496-zdu
in repository https://gitbox.apache.org/repos/asf/ozone.git


The following commit(s) were added to refs/heads/HDDS-14496-zdu by this push:
     new 44b89fdac5c HDDS-15718. Pass write pipeline version from client to 
Datanode on write (#10991)
44b89fdac5c is described below

commit 44b89fdac5c7dd41f5ca82c91026b5765d3e0293
Author: Zita Dombi <[email protected]>
AuthorDate: Wed Sep 9 18:32:52 2026 +0200

    HDDS-15718. Pass write pipeline version from client to Datanode on write 
(#10991)
---
 .../hdds/scm/storage/BlockDataStreamOutput.java    |   3 +-
 .../hadoop/hdds/protocol/DatanodeDetails.java      |  34 +++-
 .../apache/hadoop/hdds/scm/pipeline/Pipeline.java  |  15 ++
 .../hdds/scm/storage/ContainerProtocolCalls.java   |  12 +-
 .../hadoop/hdds/scm/utils/ClientCommandsUtils.java |  43 +++++
 .../scm/storage/TestContainerProtocolCalls.java    | 186 +++++++++++++++++++++
 .../hdds/scm/utils/TestClientCommandsUtils.java    |  77 +++++++++
 .../server/ratis/ContainerStateMachine.java        |   4 +
 .../transport/server/ratis/DispatcherContext.java  |  26 +++
 .../ozone/container/keyvalue/KeyValueHandler.java  |   9 +-
 .../server/ratis/TestDispatcherContext.java        |  51 ++++++
 .../container/keyvalue/TestKeyValueHandler.java    |  84 ++++++++++
 .../src/main/proto/DatanodeClientProtocol.proto    |   3 +
 13 files changed, 538 insertions(+), 9 deletions(-)

diff --git 
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockDataStreamOutput.java
 
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockDataStreamOutput.java
index e3fff7528d4..cdf353aff3f 100644
--- 
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockDataStreamOutput.java
+++ 
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockDataStreamOutput.java
@@ -218,7 +218,8 @@ private DataStreamOutput setupStream(Pipeline pipeline) 
throws IOException {
         ContainerProtos.ContainerCommandRequestProto.newBuilder()
             .setCmdType(streamInitType)
             .setContainerID(blockID.get().getContainerID())
-            .setDatanodeUuid(id).setWriteChunk(writeChunkRequest);
+            .setDatanodeUuid(id).setWriteChunk(writeChunkRequest)
+            .setWritePipelineVersion(pipeline.getWriteVersion());
 
     if (tokenString != null) {
       builder.setEncodedToken(tokenString);
diff --git 
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
 
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
index c354f7f4d93..befb0ed373b 100644
--- 
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
+++ 
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
@@ -91,6 +91,10 @@ public class DatanodeDetails extends NodeImpl implements 
Comparable<DatanodeDeta
   private volatile long persistedOpStateExpiryEpochSec;
   private HDDSVersion initialVersion;
   private volatile HDDSVersion currentVersion;
+  // The raw serialized currentVersion as it arrived off the wire. Kept 
separate from
+  // currentVersion so that a reader which deserializes an unknown newer 
version to
+  // UNKNOWN_VERSION still forwards the original value to datanodes instead of 
-1.
+  private volatile int serializedCurrentVersion;
 
   private DatanodeDetails(Builder b) {
     super(b.hostName, b.networkLocation, NetConstants.NODE_COST_DEFAULT);
@@ -107,6 +111,7 @@ private DatanodeDetails(Builder b) {
     persistedOpStateExpiryEpochSec = b.persistedOpStateExpiryEpochSec;
     initialVersion = b.initialVersion;
     currentVersion = b.currentVersion;
+    serializedCurrentVersion = b.serializedCurrentVersion;
     if (b.networkName != null) {
       setNetworkName(b.networkName);
     }
@@ -135,6 +140,7 @@ public DatanodeDetails(DatanodeDetails datanodeDetails) {
         datanodeDetails.getPersistedOpStateExpiryEpochSec();
     this.initialVersion = datanodeDetails.getInitialVersion();
     this.currentVersion = datanodeDetails.getCurrentVersion();
+    this.serializedCurrentVersion = 
datanodeDetails.getSerializedCurrentVersion();
   }
 
   public static Codec<DatanodeDetails> getCodec() {
@@ -476,7 +482,7 @@ public static DatanodeDetails.Builder newBuilder(
           datanodeDetailsProto.getPersistedOpStateExpiry());
     }
     if (datanodeDetailsProto.hasCurrentVersion()) {
-      
builder.setCurrentVersion(HDDSVersion.deserialize(datanodeDetailsProto.getCurrentVersion()));
+      builder.setCurrentVersion(datanodeDetailsProto.getCurrentVersion());
     }
     return builder;
   }
@@ -613,7 +619,7 @@ public HddsProtos.DatanodeDetailsProto.Builder 
toProtoBuilder(
       }
     }
 
-    builder.setCurrentVersion(versionOverride != null ? 
versionOverride.serialize() : currentVersion.serialize());
+    builder.setCurrentVersion(versionOverride != null ? 
versionOverride.serialize() : serializedCurrentVersion);
 
     return builder;
   }
@@ -660,8 +666,19 @@ public HDDSVersion getCurrentVersion() {
     return currentVersion;
   }
 
+  /**
+   * @return the raw serialized currentVersion as it arrived off the wire. 
Unlike
+   *     {@link #getCurrentVersion()} this preserves the original value even 
when it cannot be
+   *     deserialized by this component. Use this when forwarding the version 
to another
+   *     component, e.g. the write pipeline version sent to datanodes.
+   */
+  public int getSerializedCurrentVersion() {
+    return serializedCurrentVersion;
+  }
+
   public void setCurrentVersion(HDDSVersion currentVersion) {
     this.currentVersion = currentVersion;
+    this.serializedCurrentVersion = currentVersion.serialize();
   }
 
   @Override
@@ -746,6 +763,7 @@ public static final class Builder {
     private long persistedOpStateExpiryEpochSec = 0;
     private HDDSVersion initialVersion = HDDSVersion.DEFAULT_VERSION;
     private HDDSVersion currentVersion = HDDSVersion.DEFAULT_VERSION;
+    private int serializedCurrentVersion = 
HDDSVersion.DEFAULT_VERSION.serialize();
 
     /**
      * Default private constructor. To create Builder instance use
@@ -969,6 +987,18 @@ public Builder setInitialVersion(HDDSVersion v) {
 
     public Builder setCurrentVersion(HDDSVersion v) {
       this.currentVersion = v;
+      this.serializedCurrentVersion = v.serialize();
+      return this;
+    }
+
+    /**
+     * Sets the currentVersion from its raw serialized value, preserving the 
original even if it
+     * cannot be deserialized (in which case {@link #getCurrentVersion()} is
+     * {@link HDDSVersion#UNKNOWN_VERSION} while {@link 
#getSerializedCurrentVersion()} keeps it).
+     */
+    public Builder setCurrentVersion(int serialized) {
+      this.serializedCurrentVersion = serialized;
+      this.currentVersion = HDDSVersion.deserialize(serialized);
       return this;
     }
 
diff --git 
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
 
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
index 8c084d7c8eb..189a40de9c4 100644
--- 
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
+++ 
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
@@ -39,6 +39,7 @@
 import org.apache.commons.lang3.builder.EqualsBuilder;
 import org.apache.commons.lang3.builder.HashCodeBuilder;
 import org.apache.hadoop.hdds.ComponentVersion;
+import org.apache.hadoop.hdds.HDDSVersion;
 import org.apache.hadoop.hdds.client.ECReplicationConfig;
 import org.apache.hadoop.hdds.client.ReplicatedReplicationConfig;
 import org.apache.hadoop.hdds.client.ReplicationConfig;
@@ -289,6 +290,20 @@ public DatanodeDetails getFirstNode(Set<DatanodeDetails> 
excluded)
         "All nodes are excluded: Pipeline=%s, excluded=%s", id, excluded));
   }
 
+  /**
+   * Returns the serialized component version this pipeline should execute 
writes at. All nodes in
+   * a pipeline are expected to share the same version, so we use the first 
node's. The raw
+   * serialized value is returned (see {@link 
DatanodeDetails#getSerializedCurrentVersion()}) so a
+   * client that cannot deserialize a newer datanode version still forwards 
the original value
+   * rather than {@link HDDSVersion#UNKNOWN_VERSION}.
+   */
+  public int getWriteVersion() throws IOException {
+    if (nodeStatus.isEmpty()) {
+      throw new IOException(String.format("Pipeline=%s is empty", id));
+    }
+    return nodeStatus.keySet().iterator().next().getSerializedCurrentVersion();
+  }
+
   public DatanodeDetails getClosestNode() throws IOException {
     return getClosestNode(null);
   }
diff --git 
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java
 
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java
index d2898c395d9..d030db429ae 100644
--- 
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java
+++ 
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java
@@ -365,7 +365,8 @@ public static ContainerProtos.FinalizeBlockResponseProto 
finalizeBlock(
         
ContainerCommandRequestProto.newBuilder().setCmdType(Type.FinalizeBlock)
             .setContainerID(blockID.getContainerID())
             .setDatanodeUuid(id)
-            .setFinalizeBlock(finalizeBlockRequest);
+            .setFinalizeBlock(finalizeBlockRequest)
+            
.setWritePipelineVersion(xceiverClient.getPipeline().getWriteVersion());
     if (token != null) {
       builder.setEncodedToken(token.encodeToUrlString());
     }
@@ -396,7 +397,8 @@ public static ContainerCommandRequestProto 
getPutBlockRequest(
         ContainerCommandRequestProto.newBuilder().setCmdType(Type.PutBlock)
             .setContainerID(containerBlockData.getBlockID().getContainerID())
             .setDatanodeUuid(id)
-            .setPutBlock(createBlockRequest);
+            .setPutBlock(createBlockRequest)
+            .setWritePipelineVersion(pipeline.getWriteVersion());
     if (tokenString != null) {
       builder.setEncodedToken(tokenString);
     }
@@ -536,7 +538,8 @@ public static XceiverClientReply writeChunkAsync(
             .setCmdType(Type.WriteChunk)
             .setContainerID(blockID.getContainerID())
             .setDatanodeUuid(id)
-            .setWriteChunk(writeChunkRequest);
+            .setWriteChunk(writeChunkRequest)
+            
.setWritePipelineVersion(xceiverClient.getPipeline().getWriteVersion());
 
     if (tokenString != null) {
       builder.setEncodedToken(tokenString);
@@ -594,7 +597,8 @@ public static PutSmallFileResponseProto writeSmallFile(
             .setCmdType(Type.PutSmallFile)
             .setContainerID(blockID.getContainerID())
             .setDatanodeUuid(id)
-            .setPutSmallFile(putSmallFileRequest);
+            .setPutSmallFile(putSmallFileRequest)
+            .setWritePipelineVersion(client.getPipeline().getWriteVersion());
     if (token != null) {
       builder.setEncodedToken(token.encodeToUrlString());
     }
diff --git 
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/utils/ClientCommandsUtils.java
 
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/utils/ClientCommandsUtils.java
index 4edec24edcb..22cd3f0160b 100644
--- 
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/utils/ClientCommandsUtils.java
+++ 
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/utils/ClientCommandsUtils.java
@@ -17,13 +17,18 @@
 
 package org.apache.hadoop.hdds.scm.utils;
 
+import org.apache.hadoop.hdds.HDDSVersion;
 import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 /**
  * These methods should be merged with other similar utility classes.
  */
 public final class ClientCommandsUtils {
 
+  private static final Logger LOG = 
LoggerFactory.getLogger(ClientCommandsUtils.class);
+
   /** Utility classes should not be constructed. **/
   private ClientCommandsUtils() {
 
@@ -46,4 +51,42 @@ public static ContainerProtos.ReadChunkVersion 
getReadChunkVersion(
       return ContainerProtos.ReadChunkVersion.V0;
     }
   }
+
+  /**
+   * Returns the write pipeline (component) version the datanode should 
execute the write at,
+   * derived from the value the client forwarded from the SCM-provided 
pipeline version.
+   * {@link HDDSVersion#ZDU} is the lowest version write-path versioning can 
use, so it is the
+   * floor: any version below ZDU (including an absent field from a client 
predating zero downtime
+   * upgrade support, or a value this datanode cannot deserialize) is rounded 
up to ZDU. ZDU is the
+   * first value that is unambiguous across both the component-version and the 
(legacy)
+   * layout-feature domains, so any {@code isAllowed} comparison on a datanode 
that has not
+   * finalized ZDU is safe. No write-path versioning feature predates ZDU, so 
this loses no
+   * behavior. New clients may still forward a pre-ZDU current version (e.g. 
when HDDS is not yet
+   * finalized for ZDU); the datanode rounds it up here so there is no issue.
+   */
+  public static HDDSVersion 
getWritePipelineVersion(ContainerProtos.ContainerCommandRequestProto request) {
+    final HDDSVersion defaultWriteVersion = HDDSVersion.ZDU;
+    // If no write version is provided, fall back to the default.
+    if (!request.hasWritePipelineVersion()) {
+      return defaultWriteVersion;
+    }
+
+    // If the write version does not deserialize to any known version, fall 
back to the default.
+    int serializedVersion = request.getWritePipelineVersion();
+    HDDSVersion writeVersion = HDDSVersion.deserialize(serializedVersion);
+    if (writeVersion == HDDSVersion.UNKNOWN_VERSION) {
+      // Should not normally happen: the version originates from SCM's view of 
the datanodes.
+      LOG.error("Datanode was given an unrecognized write pipeline version {}; 
using {} instead.",
+          serializedVersion, defaultWriteVersion);
+      return defaultWriteVersion;
+    }
+
+    // ZDU is the first version that supports write path versioning.
+    // A lower version should be replaced with the default.
+    if (!defaultWriteVersion.isSupportedBy(writeVersion)) {
+      return defaultWriteVersion;
+    }
+
+    return writeVersion;
+  }
 }
diff --git 
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerProtocolCalls.java
 
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerProtocolCalls.java
new file mode 100644
index 00000000000..8d934c24e80
--- /dev/null
+++ 
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerProtocolCalls.java
@@ -0,0 +1,186 @@
+/*
+ * 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.hadoop.hdds.scm.storage;
+
+import static 
org.apache.hadoop.hdds.protocol.MockDatanodeDetails.randomDatanodeDetails;
+import static 
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyList;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.List;
+import org.apache.hadoop.hdds.HDDSVersion;
+import org.apache.hadoop.hdds.client.BlockID;
+import org.apache.hadoop.hdds.client.RatisReplicationConfig;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import 
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.BlockData;
+import 
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChecksumType;
+import 
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfo;
+import 
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto;
+import 
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandResponseProto;
+import 
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.DatanodeBlockID;
+import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result;
+import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Type;
+import org.apache.hadoop.hdds.scm.XceiverClientReply;
+import org.apache.hadoop.hdds.scm.XceiverClientSpi;
+import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
+import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
+import org.apache.hadoop.ozone.common.Checksum;
+import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+
+/**
+ * Test that {@link ContainerProtocolCalls} stamps the pipeline write version 
on
+ * outgoing write requests.
+ */
+public class TestContainerProtocolCalls {
+
+  private static final int VERSION = HDDSVersion.ZDU.serialize();
+
+  private static Pipeline pipelineWithVersion(int version) {
+    List<DatanodeDetails> nodes = Arrays.asList(
+        randomDatanodeDetails(), randomDatanodeDetails(), 
randomDatanodeDetails());
+    nodes.forEach(dn -> 
dn.setCurrentVersion(HDDSVersion.deserialize(version)));
+    return Pipeline.newBuilder()
+        .setId(PipelineID.randomId())
+        .setState(Pipeline.PipelineState.OPEN)
+        .setReplicationConfig(RatisReplicationConfig.getInstance(THREE))
+        .setNodes(nodes)
+        .build();
+  }
+
+  private static XceiverClientSpi clientWith(Pipeline pipeline) {
+    XceiverClientSpi client = mock(XceiverClientSpi.class);
+    when(client.getPipeline()).thenReturn(pipeline);
+    return client;
+  }
+
+  private static ChunkInfo chunkInfo(byte[] data) throws IOException {
+    return ChunkInfo.newBuilder()
+        .setChunkName("chunk")
+        .setOffset(0)
+        .setLen(data.length)
+        .setChecksumData(new Checksum(ChecksumType.CRC32, 256)
+            .computeChecksum(data).getProtoBufMessage())
+        .build();
+  }
+
+  @Test
+  void putBlockRequestCarriesPipelineWriteVersion() throws IOException {
+    Pipeline pipeline = pipelineWithVersion(VERSION);
+    BlockData blockData = BlockData.newBuilder()
+        .setBlockID(DatanodeBlockID.newBuilder()
+            .setContainerID(1).setLocalID(1).build())
+        .build();
+
+    ContainerCommandRequestProto request =
+        ContainerProtocolCalls.getPutBlockRequest(pipeline, blockData, true, 
null);
+
+    assertEquals(VERSION, request.getWritePipelineVersion());
+  }
+
+  @Test
+  void writeChunkRequestCarriesPipelineWriteVersion() throws Exception {
+    Pipeline pipeline = pipelineWithVersion(VERSION);
+    XceiverClientSpi client = clientWith(pipeline);
+    
when(client.sendCommandAsync(any())).thenReturn(mock(XceiverClientReply.class));
+    byte[] data = "data".getBytes(StandardCharsets.UTF_8);
+
+    ContainerProtocolCalls.writeChunkAsync(client, chunkInfo(data), new 
BlockID(1, 1),
+        ByteString.copyFrom(data), null, 0, null, false, true);
+
+    ArgumentCaptor<ContainerCommandRequestProto> captor =
+        ArgumentCaptor.forClass(ContainerCommandRequestProto.class);
+    verify(client).sendCommandAsync(captor.capture());
+    assertEquals(VERSION, captor.getValue().getWritePipelineVersion());
+  }
+
+  @Test
+  void writeSmallFileRequestCarriesPipelineWriteVersion() throws IOException {
+    Pipeline pipeline = pipelineWithVersion(VERSION);
+    XceiverClientSpi client = clientWith(pipeline);
+    when(client.sendCommand(any(), anyList())).thenReturn(
+        response(Type.PutSmallFile));
+
+    ContainerProtocolCalls.writeSmallFile(client, new BlockID(1, 1),
+        "data".getBytes(StandardCharsets.UTF_8), null);
+
+    assertEquals(VERSION, sentRequest(client).getWritePipelineVersion());
+  }
+
+  @Test
+  void finalizeBlockRequestCarriesPipelineWriteVersion() throws IOException {
+    Pipeline pipeline = pipelineWithVersion(VERSION);
+    XceiverClientSpi client = clientWith(pipeline);
+    when(client.sendCommand(any(), anyList())).thenReturn(
+        response(Type.FinalizeBlock));
+
+    ContainerProtocolCalls.finalizeBlock(client,
+        DatanodeBlockID.newBuilder().setContainerID(1).setLocalID(1).build(), 
null);
+
+    assertEquals(VERSION, sentRequest(client).getWritePipelineVersion());
+  }
+
+  @Test
+  void writeVersionPreservesOriginalSerializedValueFromWire() throws 
IOException {
+    // A datanode current version this (older) client cannot deserialize: it 
becomes
+    // UNKNOWN_VERSION locally, but the original serialized int must still be 
forwarded.
+    DatanodeDetails dn = DatanodeDetails.getFromProtoBuf(
+        randomDatanodeDetails().getProtoBufMessage().toBuilder()
+            .setCurrentVersion(Integer.MAX_VALUE).build());
+    assertEquals(HDDSVersion.UNKNOWN_VERSION, dn.getCurrentVersion());
+
+    Pipeline pipeline = Pipeline.newBuilder()
+        .setId(PipelineID.randomId())
+        .setState(Pipeline.PipelineState.OPEN)
+        .setReplicationConfig(RatisReplicationConfig.getInstance(THREE))
+        .setNodes(Arrays.asList(dn, randomDatanodeDetails(), 
randomDatanodeDetails()))
+        .build();
+    BlockData blockData = BlockData.newBuilder()
+        .setBlockID(DatanodeBlockID.newBuilder()
+            .setContainerID(1).setLocalID(1).build())
+        .build();
+
+    ContainerCommandRequestProto request =
+        ContainerProtocolCalls.getPutBlockRequest(pipeline, blockData, true, 
null);
+
+    assertEquals(Integer.MAX_VALUE, request.getWritePipelineVersion());
+  }
+
+  private static ContainerCommandResponseProto response(Type type) {
+    return ContainerCommandResponseProto.newBuilder()
+        .setCmdType(type)
+        .setResult(Result.SUCCESS)
+        .build();
+  }
+
+  private static ContainerCommandRequestProto sentRequest(XceiverClientSpi 
client)
+      throws IOException {
+    ArgumentCaptor<ContainerCommandRequestProto> captor =
+        ArgumentCaptor.forClass(ContainerCommandRequestProto.class);
+    verify(client).sendCommand(captor.capture(), anyList());
+    return captor.getValue();
+  }
+}
diff --git 
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/utils/TestClientCommandsUtils.java
 
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/utils/TestClientCommandsUtils.java
new file mode 100644
index 00000000000..b37c86085bc
--- /dev/null
+++ 
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/utils/TestClientCommandsUtils.java
@@ -0,0 +1,77 @@
+/*
+ * 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.hadoop.hdds.scm.utils;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import org.apache.hadoop.hdds.HDDSVersion;
+import 
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto;
+import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Type;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Test for {@link ClientCommandsUtils}.
+ */
+public class TestClientCommandsUtils {
+
+  private static ContainerCommandRequestProto.Builder baseRequest() {
+    return ContainerCommandRequestProto.newBuilder()
+        .setCmdType(Type.WriteChunk)
+        .setContainerID(1)
+        .setDatanodeUuid("dn");
+  }
+
+  @Test
+  void returnsWritePipelineVersionWhenAtLeastZdu() {
+    ContainerCommandRequestProto request = baseRequest()
+        .setWritePipelineVersion(HDDSVersion.ZDU.serialize())
+        .build();
+
+    assertEquals(HDDSVersion.ZDU,
+        ClientCommandsUtils.getWritePipelineVersion(request));
+  }
+
+  @Test
+  void roundsUpToZduWhenBelowZdu() {
+    // A recognized pre-ZDU version must be rounded up to ZDU (the 
write-versioning floor).
+    ContainerCommandRequestProto request = baseRequest()
+        .setWritePipelineVersion(HDDSVersion.SHORT_CIRCUIT_READS.serialize())
+        .build();
+
+    assertEquals(HDDSVersion.ZDU,
+        ClientCommandsUtils.getWritePipelineVersion(request));
+  }
+
+  @Test
+  void defaultsToZduWhenAbsent() {
+    ContainerCommandRequestProto request = baseRequest().build();
+
+    assertEquals(HDDSVersion.ZDU,
+        ClientCommandsUtils.getWritePipelineVersion(request));
+  }
+
+  @Test
+  void fallsBackToZduForUnknownVersion() {
+    ContainerCommandRequestProto request = baseRequest()
+        .setWritePipelineVersion(HDDSVersion.UNKNOWN_VERSION.serialize())
+        .build();
+
+    assertEquals(HDDSVersion.ZDU,
+        ClientCommandsUtils.getWritePipelineVersion(request));
+  }
+}
diff --git 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java
 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java
index 9a3e988e27a..5db6985b1f2 100644
--- 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java
+++ 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java
@@ -68,6 +68,7 @@
 import org.apache.hadoop.hdds.scm.ScmConfigKeys;
 import 
org.apache.hadoop.hdds.scm.container.common.helpers.ContainerNotOpenException;
 import 
org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException;
+import org.apache.hadoop.hdds.scm.utils.ClientCommandsUtils;
 import org.apache.hadoop.hdds.utils.Cache;
 import org.apache.hadoop.hdds.utils.ResourceCache;
 import org.apache.hadoop.ozone.HddsDatanodeService;
@@ -605,6 +606,7 @@ private CompletableFuture<Message> writeStateMachineData(
             .setLogIndex(entryIndex)
             .setStage(DispatcherContext.WriteChunkStage.WRITE_DATA)
             .setContainer2BCSIDMap(container2BCSIDMap)
+            
.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(requestProto))
             .build();
     CompletableFuture<Message> raftFuture = new CompletableFuture<>();
     // ensure the write chunk happens asynchronously in writeChunkExecutor 
pool thread.
@@ -751,6 +753,7 @@ public CompletableFuture<DataStream> 
stream(RaftClientRequest request) {
                 .newBuilder(DispatcherContext.Op.STREAM_INIT)
                 .setStage(DispatcherContext.WriteChunkStage.WRITE_DATA)
                 .setContainer2BCSIDMap(container2BCSIDMap)
+                
.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(requestProto))
                 .build();
         final DataChannel channel = getStreamDataChannel(requestProto, 
context);
         final ExecutorService chunkExecutor = requestProto.hasWriteChunk() ?
@@ -1150,6 +1153,7 @@ public CompletableFuture<Message> 
applyTransaction(TransactionContext trx) {
       Objects.requireNonNull(context, "context == null");
       final ContainerCommandRequestProto requestProto = context.getLogProto();
       final Type cmdType = requestProto.getCmdType();
+      
builder.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(requestProto));
       // Make sure that in write chunk, the user data is not set
       if (cmdType == Type.WriteChunk) {
         Preconditions
diff --git 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/DispatcherContext.java
 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/DispatcherContext.java
index e0f0ccbe193..0bf3221fde0 100644
--- 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/DispatcherContext.java
+++ 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/DispatcherContext.java
@@ -19,6 +19,7 @@
 
 import java.util.Map;
 import java.util.Objects;
+import org.apache.hadoop.hdds.HDDSVersion;
 import org.apache.hadoop.hdds.annotation.InterfaceAudience;
 import org.apache.hadoop.hdds.annotation.InterfaceStability;
 import org.apache.hadoop.util.Time;
@@ -53,6 +54,10 @@ public final class DispatcherContext {
 
   private final Map<Long, Long> container2BCSIDMap;
 
+  // the component version the client asked the datanode to execute this write 
at
+  // defaults to HDDSVersion.ZDU (the safe cross-domain floor) when the client 
did not supply one
+  private final HDDSVersion writeVersion;
+
   private final boolean releaseSupported;
   private volatile Runnable releaseMethod;
 
@@ -137,6 +142,7 @@ private DispatcherContext(Builder b) {
     this.logIndex = b.logIndex;
     this.stage = b.stage;
     this.container2BCSIDMap = b.container2BCSIDMap;
+    this.writeVersion = b.writeVersion;
     this.releaseSupported = b.releaseSupported;
   }
 
@@ -161,6 +167,14 @@ public Map<Long, Long> getContainer2BCSIDMap() {
     return container2BCSIDMap;
   }
 
+  /**
+   * @return the component version the client requested this write be executed 
at.
+   *     Defaults to {@link HDDSVersion#ZDU}, the safe cross-domain floor, 
when the client did not supply one.
+   */
+  public HDDSVersion getWriteVersion() {
+    return writeVersion;
+  }
+
   public long getStartTime() {
     return startTime;
   }
@@ -198,6 +212,7 @@ public static final class Builder {
     private long term;
     private long logIndex;
     private Map<Long, Long> container2BCSIDMap;
+    private HDDSVersion writeVersion = HDDSVersion.ZDU;
     private boolean releaseSupported;
 
     private Builder(Op op) {
@@ -248,6 +263,17 @@ public Builder setContainer2BCSIDMap(Map<Long, Long> map) {
       return this;
     }
 
+    /**
+     * Sets the component version the client requested this write be executed 
at.
+     *
+     * @param version the write pipeline (component) version
+     * @return Builder
+     */
+    public Builder setWriteVersion(HDDSVersion version) {
+      this.writeVersion = version;
+      return this;
+    }
+
     public Builder setReleaseSupported(boolean releaseSupported) {
       this.releaseSupported = releaseSupported;
       return this;
diff --git 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java
 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java
index a32a56a0267..d361b2f5c61 100644
--- 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java
+++ 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java
@@ -123,6 +123,7 @@
 import org.apache.hadoop.hdds.scm.storage.BlockInputStream;
 import org.apache.hadoop.hdds.scm.storage.BlockLocationInfo;
 import org.apache.hadoop.hdds.scm.storage.ChunkInputStream;
+import org.apache.hadoop.hdds.scm.utils.ClientCommandsUtils;
 import org.apache.hadoop.hdds.security.token.OzoneBlockTokenIdentifier;
 import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
 import org.apache.hadoop.hdds.utils.FaultInjector;
@@ -1106,7 +1107,9 @@ ContainerCommandResponseProto handleWriteChunk(
 
       ChunkBuffer data = null;
       if (dispatcherContext == null) {
-        dispatcherContext = DispatcherContext.getHandleWriteChunk();
+        dispatcherContext = 
DispatcherContext.newBuilder(DispatcherContext.Op.HANDLE_WRITE_CHUNK)
+            
.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(request))
+            .build();
       }
       final boolean isWrite = dispatcherContext.getStage().isWrite();
       if (isWrite) {
@@ -1265,7 +1268,9 @@ ContainerCommandResponseProto handlePutSmallFile(
       ChunkBuffer data = ChunkBuffer.wrap(
           putSmallFileReq.getData().asReadOnlyByteBufferList());
       if (dispatcherContext == null) {
-        dispatcherContext = DispatcherContext.getHandlePutSmallFile();
+        dispatcherContext = 
DispatcherContext.newBuilder(DispatcherContext.Op.HANDLE_PUT_SMALL_FILE)
+            
.setWriteVersion(ClientCommandsUtils.getWritePipelineVersion(request))
+            .build();
       }
 
       BlockID blockID = blockData.getBlockID();
diff --git 
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/TestDispatcherContext.java
 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/TestDispatcherContext.java
new file mode 100644
index 00000000000..b8edc6de325
--- /dev/null
+++ 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/TestDispatcherContext.java
@@ -0,0 +1,51 @@
+/*
+ * 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.hadoop.ozone.container.common.transport.server.ratis;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import org.apache.hadoop.hdds.HDDSVersion;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Test for {@link DispatcherContext} write version handling (HDDS-15718).
+ */
+public class TestDispatcherContext {
+
+  @Test
+  void writeVersionDefaultsToZdu() {
+    DispatcherContext context = DispatcherContext
+        .newBuilder(DispatcherContext.Op.WRITE_STATE_MACHINE_DATA)
+        .build();
+
+    assertEquals(HDDSVersion.ZDU,
+        context.getWriteVersion());
+  }
+
+  @Test
+  void writeVersionIsCarried() {
+    // Use a non-default version so this exercises the setter, not the ZDU 
default.
+    DispatcherContext context = DispatcherContext
+        .newBuilder(DispatcherContext.Op.WRITE_STATE_MACHINE_DATA)
+        .setWriteVersion(HDDSVersion.SHORT_CIRCUIT_READS)
+        .build();
+
+    assertEquals(HDDSVersion.SHORT_CIRCUIT_READS,
+        context.getWriteVersion());
+  }
+}
diff --git 
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java
 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java
index 2670f1585a2..92cb3adad12 100644
--- 
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java
+++ 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java
@@ -56,6 +56,7 @@
 import com.google.common.collect.Sets;
 import java.io.File;
 import java.io.IOException;
+import java.lang.reflect.Field;
 import java.nio.ByteBuffer;
 import java.nio.file.Files;
 import java.nio.file.Path;
@@ -72,6 +73,7 @@
 import org.apache.commons.io.FileUtils;
 import org.apache.hadoop.conf.StorageUnit;
 import org.apache.hadoop.fs.FileUtil;
+import org.apache.hadoop.hdds.HDDSVersion;
 import org.apache.hadoop.hdds.client.BlockID;
 import org.apache.hadoop.hdds.conf.OzoneConfiguration;
 import org.apache.hadoop.hdds.protocol.DatanodeDetails;
@@ -84,6 +86,7 @@
 import 
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReplicaProto;
 import org.apache.hadoop.hdds.scm.ScmConfigKeys;
 import 
org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException;
+import org.apache.hadoop.hdds.scm.pipeline.MockPipeline;
 import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
 import org.apache.hadoop.hdds.security.token.TokenVerifier;
 import org.apache.hadoop.hdds.utils.io.RandomAccessFileChannel;
@@ -113,6 +116,7 @@
 import 
org.apache.hadoop.ozone.container.common.volume.RoundRobinVolumeChoosingPolicy;
 import org.apache.hadoop.ozone.container.common.volume.StorageVolume;
 import org.apache.hadoop.ozone.container.common.volume.VolumeSet;
+import org.apache.hadoop.ozone.container.keyvalue.interfaces.ChunkManager;
 import org.apache.hadoop.ozone.container.ozoneimpl.ContainerController;
 import 
org.apache.hadoop.ozone.container.ozoneimpl.ContainerScannerConfiguration;
 import org.apache.hadoop.ozone.container.ozoneimpl.OnDemandContainerScanner;
@@ -124,6 +128,7 @@
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
+import org.mockito.ArgumentCaptor;
 import org.mockito.Mockito;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -994,6 +999,85 @@ private HandlerWithVolumeSet createKeyValueHandler(Path 
path) throws IOException
     return new HandlerWithVolumeSet(kvHandler, volumeSet, containerSet);
   }
 
+  /**
+   * EC/standalone writes reach the datanode with a null DispatcherContext, so 
KeyValueHandler
+   * fabricates one. Verify the client-supplied write pipeline version is 
carried onto that
+   * context for WriteChunk.
+   */
+  @Test
+  public void testWriteChunkCarriesWritePipelineVersionForEc() throws 
Exception {
+    Path testDir = Files.createTempDirectory("testWriteChunkEcVersion");
+    conf.set(OZONE_SCM_CONTAINER_LAYOUT_KEY, 
ContainerLayoutVersion.FILE_PER_BLOCK.name());
+    HandlerWithVolumeSet handlerWithVolume = 
createKeyValueHandlerWithVolumeSet(testDir);
+    KeyValueHandler kvHandler = handlerWithVolume.getHandler();
+
+    long containerID = ContainerTestHelper.getTestContainerID();
+    KeyValueContainer container = createOpenContainer(containerID,
+        handlerWithVolume.getVolumeSet(), handlerWithVolume.getContainerSet());
+    ChunkManager spyChunkManager = spyChunkManager(kvHandler);
+
+    BlockID blockID = ContainerTestHelper.getTestBlockID(containerID);
+    ContainerCommandRequestProto request = ContainerTestHelper
+        .getWriteChunkRequest(MockPipeline.createSingleNodePipeline(), 
blockID, 1024)
+        .toBuilder()
+        .setWritePipelineVersion(HDDSVersion.ZDU.serialize())
+        .build();
+
+    kvHandler.handleWriteChunk(request, container, null);
+
+    ArgumentCaptor<DispatcherContext> captor = 
ArgumentCaptor.forClass(DispatcherContext.class);
+    verify(spyChunkManager).writeChunk(eq(container), any(), any(), 
any(ChunkBuffer.class), captor.capture());
+    assertEquals(HDDSVersion.ZDU, captor.getValue().getWriteVersion());
+  }
+
+  /**
+   * Same as {@link #testWriteChunkCarriesWritePipelineVersionForEc()} but for 
PutSmallFile.
+   */
+  @Test
+  public void testPutSmallFileCarriesWritePipelineVersionForEc() throws 
Exception {
+    Path testDir = Files.createTempDirectory("testPutSmallFileEcVersion");
+    conf.set(OZONE_SCM_CONTAINER_LAYOUT_KEY, 
ContainerLayoutVersion.FILE_PER_BLOCK.name());
+    HandlerWithVolumeSet handlerWithVolume = 
createKeyValueHandlerWithVolumeSet(testDir);
+    KeyValueHandler kvHandler = handlerWithVolume.getHandler();
+
+    long containerID = ContainerTestHelper.getTestContainerID();
+    KeyValueContainer container = createOpenContainer(containerID,
+        handlerWithVolume.getVolumeSet(), handlerWithVolume.getContainerSet());
+    ChunkManager spyChunkManager = spyChunkManager(kvHandler);
+
+    BlockID blockID = ContainerTestHelper.getTestBlockID(containerID);
+    ContainerCommandRequestProto request = ContainerTestHelper
+        .getWriteSmallFileRequest(MockPipeline.createSingleNodePipeline(), 
blockID, 1024)
+        .toBuilder()
+        .setWritePipelineVersion(HDDSVersion.ZDU.serialize())
+        .build();
+
+    kvHandler.handlePutSmallFile(request, container, null);
+
+    ArgumentCaptor<DispatcherContext> captor = 
ArgumentCaptor.forClass(DispatcherContext.class);
+    verify(spyChunkManager).writeChunk(eq(container), any(), any(), 
any(ChunkBuffer.class), captor.capture());
+    assertEquals(HDDSVersion.ZDU, captor.getValue().getWriteVersion());
+  }
+
+  private static ChunkManager spyChunkManager(KeyValueHandler kvHandler) 
throws Exception {
+    Field field = KeyValueHandler.class.getDeclaredField("chunkManager");
+    field.setAccessible(true);
+    ChunkManager spy = spy((ChunkManager) field.get(kvHandler));
+    field.set(kvHandler, spy);
+    return spy;
+  }
+
+  private KeyValueContainer createOpenContainer(long containerID,
+      MutableVolumeSet volumeSet, ContainerSet containerSet) throws 
IOException {
+    KeyValueContainerData containerData = new KeyValueContainerData(
+        containerID, ContainerLayoutVersion.FILE_PER_BLOCK,
+        (long) StorageUnit.GB.toBytes(1), UUID.randomUUID().toString(), 
DATANODE_UUID);
+    KeyValueContainer container = new KeyValueContainer(containerData, conf);
+    container.create(volumeSet, new RoundRobinVolumeChoosingPolicy(), 
CLUSTER_ID);
+    containerSet.addContainer(container);
+    return container;
+  }
+
   private static class HandlerWithVolumeSet {
     private final KeyValueHandler handler;
     private final MutableVolumeSet volumeSet;
diff --git 
a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto 
b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto
index 72d50c975a5..385cdb319cf 100644
--- a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto
+++ b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto
@@ -203,6 +203,9 @@ message ContainerCommandRequestProto {
 
   optional   GetContainerChecksumInfoRequestProto getContainerChecksumInfo = 
27;
   optional   ReadBlockRequestProto readBlock = 28;
+  // Version that datanodes should use for this write operation.
+  // Set by SCM and forwarded by the client to the datanodes.
+  optional   uint32 writePipelineVersion = 29;
 
   // clientId and callId are used to distinguish different requests from 
different local clients for shortCircuitRead
   optional   bytes clientId = 100;


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

Reply via email to