peterxcli commented on code in PR #11257:
URL: https://github.com/apache/ozone/pull/11257#discussion_r4104776948


##########
hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OMRatisRequestContext.java:
##########


Review Comment:
   <details>
   <summary>suggest diff(gen with help of ai, but I do think this really make 
code cleaner): </summary>
   
   ```diff
   diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/execution/OMExecutionFlow.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/execution/OMExecutionFlow.java
   index bc33cb5c45..9b3599ba11 100644
   --- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/execution/OMExecutionFlow.java
   +++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/execution/OMExecutionFlow.java
   @@ -24,11 +24,11 @@
    import org.apache.hadoop.ozone.om.OMPerformanceMetrics;
    import org.apache.hadoop.ozone.om.OzoneManager;
    import org.apache.hadoop.ozone.om.helpers.OMAuditLogger;
   -import org.apache.hadoop.ozone.om.ratis.OMRatisRequestContext;
    import org.apache.hadoop.ozone.om.ratis.utils.OzoneManagerRatisUtils;
    import org.apache.hadoop.ozone.om.request.OMClientRequest;
    import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
    import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMResponse;
   +import org.apache.hadoop.ozone.security.STSSecurityUtil;
    
    /**
     * entry for execution flow for write request.
   @@ -75,10 +75,7 @@ private OMResponse submitExecutionToRatis(OMRequest 
request, boolean isWrite) th
          }
        } else {
          try {
   -        // We capture the ThreadLocal context in OzoneManager so that it 
will be available
   -        // in a separate read thread in OzoneManagerStateMachine#query.
   -        // Write requests already implemented this logic in the preExecute.
   -        requestToSubmit = OMRatisRequestContext.captureIntoRequest(request, 
ozoneManager);
   +        requestToSubmit = captureReadContext(request);
          } catch (IOException ex) {
            return OzoneManagerRatisUtils.createErrorResponse(request, ex);
          }
   @@ -91,4 +88,19 @@ private OMResponse submitExecutionToRatis(OMRequest 
request, boolean isWrite) th
        }
        return response;
      }
   +
   +  /**
   +   * Serializes the caller's thread-local context into a read request so 
that
   +   * OzoneManagerStateMachine#query can restore it, possibly on a separate 
thread.
   +   * Write requests do the same in OMClientRequest#preExecute.
   +   */
   +  private OMRequest captureReadContext(OMRequest request) throws 
IOException {
   +    OMRequest.Builder requestBuilder = request.toBuilder()
   +        .setUserInfo(OMClientRequest.getAuthenticatedUserInfo(request));
   +    if (requestBuilder.hasS3Authentication()) {
   +      requestBuilder.setS3Authentication(
   +          
STSSecurityUtil.resolveS3Authentication(requestBuilder.getS3Authentication(), 
ozoneManager));
   +    }
   +    return requestBuilder.build();
   +  }
    }
   diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OMRatisRequestContext.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OMRatisRequestContext.java
   deleted file mode 100644
   index d46370761c..0000000000
   --- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OMRatisRequestContext.java
   +++ /dev/null
   @@ -1,195 +0,0 @@
   -/*
   - * 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.om.ratis;
   -
   -import java.io.IOException;
   -import java.net.InetAddress;
   -import java.net.UnknownHostException;
   -import java.util.Objects;
   -import org.apache.hadoop.ipc_.RPC;
   -import org.apache.hadoop.ipc_.RpcConstants;
   -import org.apache.hadoop.ipc_.Server;
   -import org.apache.hadoop.ozone.om.OzoneManager;
   -import org.apache.hadoop.ozone.om.lock.OMLockDetailsUtil;
   -import org.apache.hadoop.ozone.om.request.OMClientRequest;
   -import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
   -import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
   -import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMResponse;
   -import org.apache.hadoop.ozone.security.S3AuthenticationContext;
   -import org.apache.hadoop.security.UserGroupInformation;
   -import org.slf4j.Logger;
   -import org.slf4j.LoggerFactory;
   -
   -/**
   - * Captures request context before Ratis submission and applies it while an 
OM
   - * request is executed by Ratis.
   - *
   - * <p>Request-scoped thread-local state that is needed during Ratis queries
   - * must be serialized by {@link #captureIntoRequest(OMRequest, 
OzoneManager)}.
   - * Write requests capture their context during {@link 
OMClientRequest#preExecute(OzoneManager)}.
   - * Context used during either Ratis execution path must be applied by this 
class.
   - *
   - * <p>Ratis may execute a read either on the submitting RPC thread or on a 
separate thread. In both cases, the context
   - * already present on the execution thread is saved, replaced with the 
context reconstructed from the request, and
   - * restored by {@link #close()}. This prevents an inline read from clearing 
the original RPC context while also leaving
   - * a separate execution thread in its previous state.
   - *
   - * <p>Writes are applied on the separate StateMachineUpdater thread. They 
do not install a synthetic Hadoop RPC call,
   - * and {@link #close()} clears their context instead of restoring it. If 
Ratis starts applying writes on the submitting
   - * thread, the write path must also preserve and restore the existing 
context.
   - */
   -public final class OMRatisRequestContext implements AutoCloseable {
   -  private static final Logger LOG = 
LoggerFactory.getLogger(OMRatisRequestContext.class);
   -
   -  private final Server.Call currentCall;
   -  // Read scopes save these values only to restore the execution thread in 
close().
   -  private final Server.Call previousCall;
   -  private final S3AuthenticationContext previousS3Context;
   -  private final Operation operation;
   -
   -  private enum Operation {
   -    READ,
   -    WRITE
   -  }
   -
   -  private OMRatisRequestContext(OMRequest request, OzoneManager 
ozoneManager, Operation operation)
   -      throws IOException {
   -    Objects.requireNonNull(ozoneManager, "ozoneManager");
   -    this.operation = Objects.requireNonNull(operation, "operation");
   -    boolean isRead = operation == Operation.READ;
   -    previousCall = isRead ? Server.getCurCall().get() : null;
   -    previousS3Context = isRead ? S3AuthenticationContext.capture() : null;
   -    currentCall = isRead ? createCall(request) : null;
   -
   -    clear();
   -    try {
   -      if (currentCall != null) {
   -        Server.getCurCall().set(currentCall);
   -      }
   -      // This request-derived S3/STS context is active during both read and 
write execution.
   -      S3AuthenticationContext.fromRequest(request, 
ozoneManager.isSecurityEnabled()).applyToCurrentThread();
   -    } catch (IOException | RuntimeException ex) {
   -      close();
   -      throw ex;
   -    }
   -  }
   -
   -  /**
   -   * Captures authenticated request context on the RPC thread before 
submitting
   -   * a read request to Ratis.
   -   */
   -  public static OMRequest captureIntoRequest(OMRequest request, 
OzoneManager ozoneManager) throws IOException {
   -    Objects.requireNonNull(ozoneManager, "ozoneManager");
   -    OMRequest.Builder requestBuilder = request.toBuilder()
   -        .setUserInfo(OMClientRequest.getAuthenticatedUserInfo(request));
   -    S3AuthenticationContext.captureInto(requestBuilder, ozoneManager);
   -    return requestBuilder.build();
   -  }
   -
   -  /**
   -   * Applies context for a Ratis read, including a synthetic Hadoop RPC call
   -   * used by existing read implementations and lock accounting. This 
supports
   -   * execution on either the submitting thread or a separate Ratis thread by
   -   * restoring the execution thread's previous context when this scope 
closes.
   -   */
   -  public static OMRatisRequestContext openForRead(OMRequest request, 
OzoneManager ozoneManager)
   -      throws IOException {
   -    return new OMRatisRequestContext(request, ozoneManager, Operation.READ);
   -  }
   -
   -  /**
   -   * Applies context for a Ratis write without a Hadoop RPC call, preserving
   -   * write lock accounting through ResourceLockTracker. Writes are expected 
to
   -   * execute on the separate StateMachineUpdater thread, whose context is
   -   * cleared when this scope closes.
   -   */
   -  public static OMRatisRequestContext openForWrite(OMRequest request, 
OzoneManager ozoneManager)
   -      throws IOException {
   -    return new OMRatisRequestContext(request, ozoneManager, 
Operation.WRITE);
   -  }
   -
   -  /**
   -   * Adds query lock timings from the synthetic Hadoop RPC call to the 
response.
   -   */
   -  public OMResponse addLockDetails(OMResponse response) {
   -    return currentCall == null ? response
   -        : OMLockDetailsUtil.addToResponse(response, 
currentCall.getProcessingDetails());
   -  }
   -
   -  private static Server.Call createCall(OMRequest request) {
   -    if (!request.hasUserInfo()) {
   -      return null;
   -    }
   -
   -    OzoneManagerProtocolProtos.UserInfo userInfo = request.getUserInfo();
   -    UserGroupInformation remoteUser = userInfo.hasUserName() && 
!userInfo.getUserName().isEmpty()
   -        ? UserGroupInformation.createRemoteUser(userInfo.getUserName()) : 
null;
   -    InetAddress remoteAddress = createRemoteAddress(userInfo);
   -    if (remoteUser == null && remoteAddress == null) {
   -      return null;
   -    }
   -
   -    return new Server.Call(RpcConstants.INVALID_CALL_ID, 
RpcConstants.INVALID_RETRY_COUNT,
   -        null, null, RPC.RpcKind.RPC_PROTOCOL_BUFFER, 
RpcConstants.DUMMY_CLIENT_ID) {
   -      @Override
   -      public UserGroupInformation getRemoteUser() {
   -        return remoteUser;
   -      }
   -
   -      @Override
   -      public InetAddress getHostInetAddress() {
   -        return remoteAddress;
   -      }
   -    };
   -  }
   -
   -  private static InetAddress 
createRemoteAddress(OzoneManagerProtocolProtos.UserInfo userInfo) {
   -    if (!userInfo.hasRemoteAddress() || 
userInfo.getRemoteAddress().isEmpty()) {
   -      return null;
   -    }
   -
   -    try {
   -      InetAddress address = 
InetAddress.getByName(userInfo.getRemoteAddress());
   -      return userInfo.hasHostName() && !userInfo.getHostName().isEmpty()
   -          ? InetAddress.getByAddress(userInfo.getHostName(), 
address.getAddress()) : address;
   -    } catch (UnknownHostException ex) {
   -      LOG.debug("Unable to restore remote address from Ratis request 
context: {}",
   -          userInfo.getRemoteAddress(), ex);
   -      return null;
   -    }
   -  }
   -
   -  private static void clear() {
   -    Server.getCurCall().remove();
   -    S3AuthenticationContext.clear();
   -  }
   -
   -  @Override
   -  public void close() {
   -    if (operation == Operation.READ) {
   -      if (previousCall == null) {
   -        Server.getCurCall().remove();
   -      } else {
   -        Server.getCurCall().set(previousCall);
   -      }
   -      previousS3Context.applyToCurrentThread();
   -    } else {
   -      clear();
   -    }
   -  }
   -}
   diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java
   index f238bccd48..cdd86e10b5 100644
   --- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java
   +++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java
   @@ -695,7 +695,7 @@ public void close() {
       */
      @VisibleForTesting
      OMResponse runCommand(OMRequest request, TermIndex termIndex) {
   -    try (OMRatisRequestContext ignored = 
OMRatisRequestContext.openForWrite(request, ozoneManager)) {
   +    try (OMThreadContext.Scope ignored = OMThreadContext.forWrite(request, 
ozoneManager).install()) {
          ExecutionContext context = ExecutionContext.of(termIndex.getIndex(), 
termIndex);
          final OMClientResponse omClientResponse = handler.handleWriteRequest(
              request, context, ozoneManagerDoubleBuffer);
   @@ -752,9 +752,9 @@ public void loadSnapshotInfoFromDB() throws IOException {
       * @return response from OM
       */
      private Message queryCommand(OMRequest request) throws IOException {
   -    try (OMRatisRequestContext context = 
OMRatisRequestContext.openForRead(request, ozoneManager)) {
   +    try (OMThreadContext.Scope scope = OMThreadContext.forRead(request, 
ozoneManager).install()) {
          OMResponse response = handler.handleReadRequest(request);
   -      return 
OMRatisHelper.convertResponseToMessage(context.addLockDetails(response));
   +      return 
OMRatisHelper.convertResponseToMessage(scope.addLockDetails(response));
        }
      }
    
   diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/OMClientRequest.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/OMClientRequest.java
   index a6420ec296..5394a693b6 100644
   --- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/OMClientRequest.java
   +++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/OMClientRequest.java
   @@ -55,7 +55,7 @@
    import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.LayoutVersion;
    import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
    import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMResponse;
   -import org.apache.hadoop.ozone.security.S3AuthenticationContext;
   +import org.apache.hadoop.ozone.security.STSSecurityUtil;
    import org.apache.hadoop.ozone.security.acl.IAccessAuthorizer;
    import org.apache.hadoop.ozone.security.acl.OzoneObj;
    import org.apache.hadoop.ozone.security.acl.OzoneObjInfo;
   @@ -128,7 +128,8 @@ public OMRequest preExecute(OzoneManager ozoneManager)
        }
    
        if (requestBuilder.hasS3Authentication()) {
   -      S3AuthenticationContext.captureInto(requestBuilder, ozoneManager);
   +      requestBuilder.setS3Authentication(
   +          
STSSecurityUtil.resolveS3Authentication(requestBuilder.getS3Authentication(), 
ozoneManager));
        }
    
        omRequest = requestBuilder.build();
   diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/security/S3AuthenticationContext.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/security/S3AuthenticationContext.java
   deleted file mode 100644
   index e98acaf705..0000000000
   --- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/security/S3AuthenticationContext.java
   +++ /dev/null
   @@ -1,75 +0,0 @@
   -/*
   - * 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.security;
   -
   -import java.io.IOException;
   -import org.apache.hadoop.ozone.om.OzoneManager;
   -import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
   -import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.S3Authentication;
   -
   -/**
   - * Captures and applies S3 authentication state associated with the current
   - * thread.
   - */
   -public final class S3AuthenticationContext {
   -  private static final S3AuthenticationContext EMPTY = new 
S3AuthenticationContext(null, null);
   -
   -  private final S3Authentication s3Authentication;
   -  private final STSTokenIdentifier stsTokenIdentifier;
   -
   -  private S3AuthenticationContext(
   -      S3Authentication s3Authentication, STSTokenIdentifier 
stsTokenIdentifier) {
   -    this.s3Authentication = s3Authentication;
   -    this.stsTokenIdentifier = stsTokenIdentifier;
   -  }
   -
   -  public static S3AuthenticationContext capture() {
   -    return new S3AuthenticationContext(
   -        OzoneManager.getS3Auth(), OzoneManager.getStsTokenIdentifier());
   -  }
   -
   -  public static void captureInto(OMRequest.Builder requestBuilder, 
OzoneManager ozoneManager)
   -      throws IOException {
   -    if (requestBuilder.hasS3Authentication()) {
   -      requestBuilder.setS3Authentication(
   -          
STSSecurityUtil.resolveS3Authentication(requestBuilder.getS3Authentication(), 
ozoneManager));
   -    }
   -  }
   -
   -  public static S3AuthenticationContext fromRequest(
   -      OMRequest request, boolean securityEnabled) throws IOException {
   -    if (!securityEnabled || !request.hasS3Authentication()) {
   -      return EMPTY;
   -    }
   -
   -    STSSecurityUtil.ensureResolvedStsFieldsInvariants(request);
   -    S3Authentication s3Auth = request.getS3Authentication();
   -    STSTokenIdentifier stsToken = 
STSSecurityUtil.rehydrateStsTokenIdentifier(s3Auth);
   -    return new S3AuthenticationContext(s3Auth, stsToken);
   -  }
   -
   -  public void applyToCurrentThread() {
   -    OzoneManager.setS3Auth(s3Authentication);
   -    OzoneManager.setStsTokenIdentifier(stsTokenIdentifier);
   -  }
   -
   -  public static void clear() {
   -    EMPTY.applyToCurrentThread();
   -  }
   -
   -}
   diff --git 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java
 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java
   index 5ab03ba573..2857ef57ed 100644
   --- 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java
   +++ 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java
   @@ -161,7 +161,7 @@ public void 
testOzoneManagerThreadLocalsAreReviewedForRatisPropagation() {
            .collect(Collectors.toList());
    
        assertEquals(Arrays.asList("S3_AUTH", "STS_TOKEN"), threadLocalFields,
   -        "Update OMRatisRequestContext when adding request-scoped 
ThreadLocal fields to OzoneManager");
   +        "Update OMThreadContext when adding request-scoped ThreadLocal 
fields to OzoneManager");
      }
    
      // --- startTransaction tests ---
   @@ -483,7 +483,7 @@ public void testRunCommandSetsAndClearsStsThreadLocal() 
throws Exception {
      }
    
      @Test
   -  public void testRunCommandClearsStaleRequestContext() throws Exception {
   +  public void testRunCommandIsolatesAndRestoresThreadContext() throws 
Exception {
        OMRequest request = sampleWriteRequest();
        TermIndex ti = TermIndex.valueOf(1, 5);
        OMResponse expectedResponse = OMResponse.newBuilder()
   @@ -501,16 +501,19 @@ public void testRunCommandClearsStaleRequestContext() 
throws Exception {
          return clientResponse;
        });
    
   -    Server.getCurCall().set(createCall("stale-user", "stale.example.com", 
new byte[] {10, 0, 0, 1}));
   -    
OzoneManager.setS3Auth(S3Authentication.newBuilder().setAccessId("stale-access-id").build());
   -    OzoneManager.setStsTokenIdentifier(mock(STSTokenIdentifier.class));
   +    Server.Call previousCall = createCall("previous-user", 
"previous.example.com", new byte[] {10, 0, 0, 1});
   +    S3Authentication previousS3Auth = 
S3Authentication.newBuilder().setAccessId("previous-access-id").build();
   +    STSTokenIdentifier previousStsToken = mock(STSTokenIdentifier.class);
   +    Server.getCurCall().set(previousCall);
   +    OzoneManager.setS3Auth(previousS3Auth);
   +    OzoneManager.setStsTokenIdentifier(previousStsToken);
    
        OMResponse result = sm.runCommand(request, ti);
    
        assertTrue(result.getSuccess());
   -    assertNull(Server.getCurCall().get());
   -    assertNull(OzoneManager.getS3Auth());
   -    assertNull(OzoneManager.getStsTokenIdentifier());
   +    assertSame(previousCall, Server.getCurCall().get());
   +    assertSame(previousS3Auth, OzoneManager.getS3Auth());
   +    assertSame(previousStsToken, OzoneManager.getStsTokenIdentifier());
      }
    
      @Test
   ```
   </details>



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to