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

mattcasters pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git


The following commit(s) were added to refs/heads/main by this push:
     new fa2af21906 Issue #8768 : Stop the Arrow Flight server from mutating 
the shared data stream (#8769)
fa2af21906 is described below

commit fa2af21906a087ac82b5e5445270ca1bce3aba13
Author: Matt Casters <[email protected]>
AuthorDate: Tue Oct 6 11:08:47 2026 +0200

    Issue #8768 : Stop the Arrow Flight server from mutating the shared data 
stream (#8769)
    
    The producer initialized the metadata data stream from a gRPC thread. That
    object can be the client writing the batch in the same JVM, so initialize()
    replaced its allocator and flipped the writing flag. The round-trip test 
then
    saw a null id or refused to read.
    
    Clone the stream before initializing it, close Flight resources before the
    allocator, and shut the test server down with a deadline. Auth tests that
    expect a rejected login print that intent and hide Arrow's handshake stack.
---
 .../datastream/flight/ArrowFlightDataStream.java   | 27 ++++++----
 .../apache/hop/arrow/flight/ArrowFlightServer.java | 20 ++++++++
 .../apache/hop/arrow/flight/HopFlightProducer.java | 60 +++++++++++++---------
 .../ArrowFlightDataStreamAuthenticationTest.java   | 22 ++++----
 .../ArrowFlightServerAuthenticationTest.java       | 38 ++++++++------
 .../hop/arrow/flight/ExpectedFlightRejection.java  | 56 ++++++++++++++++++++
 6 files changed, 163 insertions(+), 60 deletions(-)

diff --git 
a/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/datastream/flight/ArrowFlightDataStream.java
 
b/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/datastream/flight/ArrowFlightDataStream.java
index de2d5478e6..610b64a5e3 100644
--- 
a/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/datastream/flight/ArrowFlightDataStream.java
+++ 
b/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/datastream/flight/ArrowFlightDataStream.java
@@ -333,29 +333,36 @@ public class ArrowFlightDataStream extends 
ArrowBaseDataStream {
 
   @Override
   public void close() {
-    if (readVectorSchemaRoot != null) {
-      readVectorSchemaRoot.close();
-    }
-    if (vectorSchemaRoot != null) {
-      vectorSchemaRoot.close();
-    }
-    if (rootAllocator != null) {
-      rootAllocator.close();
-    }
+    // The stream and the client still hold buffers. Close them before the 
vectors and the
+    // allocator, or BaseAllocator.close() reports the outstanding memory as a 
leak.
+    //
     if (readFlightStream != null) {
       try {
         readFlightStream.close();
       } catch (Exception e) {
         // Ignore
       }
+      readFlightStream = null;
     }
-
     if (flightClient != null) {
       try {
         flightClient.close();
       } catch (Exception e) {
         // Ignore
       }
+      flightClient = null;
+    }
+    if (readVectorSchemaRoot != null) {
+      readVectorSchemaRoot.close();
+      readVectorSchemaRoot = null;
+    }
+    if (vectorSchemaRoot != null) {
+      vectorSchemaRoot.close();
+      vectorSchemaRoot = null;
+    }
+    if (rootAllocator != null) {
+      rootAllocator.close();
+      rootAllocator = null;
     }
   }
 
diff --git 
a/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/flight/ArrowFlightServer.java
 
b/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/flight/ArrowFlightServer.java
index 981605b8c9..ebe277f38e 100644
--- 
a/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/flight/ArrowFlightServer.java
+++ 
b/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/flight/ArrowFlightServer.java
@@ -125,6 +125,9 @@ public class ArrowFlightServer {
   public void start() throws HopException {
     try {
       flightServer.start();
+      // Port 0 is replaced by the port the operating system assigned.
+      //
+      port = flightServer.getPort();
       log.logBasic(
           "Apache Arrow Flight server listening on "
               + hostname
@@ -148,4 +151,21 @@ public class ArrowFlightServer {
       throw new HopException("Unable to shut down Flight server on " + 
hostname + ":" + port, e);
     }
   }
+
+  /**
+   * Shuts the server down, waits until it has stopped, and releases its 
allocator. A caller that
+   * wants to keep the process alive uses {@link #shutdown()} and waits on the 
Flight server itself.
+   */
+  public void close() throws HopException {
+    try {
+      flightServer.close();
+      log.logBasic("Apache Arrow Flight server was shut down");
+    } catch (InterruptedException e) {
+      Thread.currentThread().interrupt();
+      throw new HopException(
+          "Interrupted while shutting down the Flight server on " + hostname + 
":" + port, e);
+    } finally {
+      allocator.close();
+    }
+  }
 }
diff --git 
a/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/flight/HopFlightProducer.java
 
b/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/flight/HopFlightProducer.java
index 30cfef7a48..806787ae7c 100644
--- 
a/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/flight/HopFlightProducer.java
+++ 
b/plugins/tech/arrow/src/main/java/org/apache/hop/arrow/flight/HopFlightProducer.java
@@ -147,38 +147,48 @@ public class HopFlightProducer extends NoOpFlightProducer 
{
             "Stream name '" + streamName + "' could not be found in the 
metadata as a data stream");
       }
       IDataStream dataStream = dataStreamMeta.getDataStream();
-      if (!(dataStream instanceof ArrowFlightDataStream flightDataStream)) {
+      if (!(dataStream instanceof ArrowFlightDataStream storedDataStream)) {
         throw new HopException(
             "Make sure to reference an Arrow Flight data stream in data stream 
element '"
                 + streamName
                 + "'.");
       }
-      flightDataStream.initialize(variables, metadataProvider, true, 
dataStreamMeta);
-      IRowMeta rowMeta = flightDataStream.buildExpectedRowMeta();
-      Schema expectedSchema = flightDataStream.buildExpectedSchema();
-
-      int bufferSize =
-          Const.toInt(
-              variables.resolve(flightDataStream.getBufferSize()),
-              ArrowFlightDataStream.DEFAULT_MAX_BUFFER_SIZE);
-      int batchSize = 
Const.toInt(variables.resolve(flightDataStream.getBatchSize()), 500);
-
-      // We use a very large queue because we don't ever want to block while 
writing.
-      // We over-size it by 5000 rows and then throw an error if we reach that.
+      // acceptPut calls this from a gRPC thread. The metadata object can be 
the instance a
+      // client in this JVM is writing with, and initialize() replaces its 
allocator and row
+      // buffer and sets the writing flag. Read the configuration from a copy.
       //
-      IRowSet rowSet = new BlockingRowSet(bufferSize + BUFFER_SIZE_OVERSHOOT);
+      ArrowFlightDataStream flightDataStream = storedDataStream.clone();
+      try {
+        flightDataStream.initialize(variables, metadataProvider, true, 
dataStreamMeta);
+        IRowMeta rowMeta = flightDataStream.buildExpectedRowMeta();
+        Schema expectedSchema = flightDataStream.buildExpectedSchema();
+
+        int bufferSize =
+            Const.toInt(
+                variables.resolve(flightDataStream.getBufferSize()),
+                ArrowFlightDataStream.DEFAULT_MAX_BUFFER_SIZE);
+        int batchSize = 
Const.toInt(variables.resolve(flightDataStream.getBatchSize()), 500);
+
+        // We use a very large queue because we don't ever want to block while 
writing.
+        // We over-size it by 5000 rows and then throw an error if we reach 
that.
+        //
+        IRowSet rowSet = new BlockingRowSet(bufferSize + 
BUFFER_SIZE_OVERSHOOT);
 
-      String hostname = 
Const.NVL(variables.resolve(flightDataStream.getHostname()), "0.0.0.0");
-      int port = Const.toInt(variables.resolve(flightDataStream.getPort()), 
33333);
-      // The endpoint we hand back needs the same scheme clients use to reach 
this server.
-      //
-      Location location =
-          flightDataStream.isTls()
-              ? Location.forGrpcTls(hostname, port)
-              : Location.forGrpcInsecure(hostname, port);
-      buffer =
-          new FlightStreamBuffer(expectedSchema, rowMeta, rowSet, bufferSize, 
batchSize, location);
-      streamMap.put(streamName, buffer);
+        String hostname = 
Const.NVL(variables.resolve(flightDataStream.getHostname()), "0.0.0.0");
+        int port = Const.toInt(variables.resolve(flightDataStream.getPort()), 
33333);
+        // The endpoint we hand back needs the same scheme clients use to 
reach this server.
+        //
+        Location location =
+            flightDataStream.isTls()
+                ? Location.forGrpcTls(hostname, port)
+                : Location.forGrpcInsecure(hostname, port);
+        buffer =
+            new FlightStreamBuffer(
+                expectedSchema, rowMeta, rowSet, bufferSize, batchSize, 
location);
+        streamMap.put(streamName, buffer);
+      } finally {
+        flightDataStream.close();
+      }
     }
     return buffer;
   }
diff --git 
a/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/datastream/flight/ArrowFlightDataStreamAuthenticationTest.java
 
b/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/datastream/flight/ArrowFlightDataStreamAuthenticationTest.java
index 71c1544b1b..3daa75ac53 100644
--- 
a/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/datastream/flight/ArrowFlightDataStreamAuthenticationTest.java
+++ 
b/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/datastream/flight/ArrowFlightDataStreamAuthenticationTest.java
@@ -26,6 +26,7 @@ import static org.mockito.Mockito.mock;
 import java.util.List;
 import org.apache.hop.arrow.flight.ArrowFlightSecurity;
 import org.apache.hop.arrow.flight.ArrowFlightServer;
+import org.apache.hop.arrow.flight.ExpectedFlightRejection;
 import org.apache.hop.core.HopClientEnvironment;
 import org.apache.hop.core.exception.HopException;
 import org.apache.hop.core.logging.ILogChannel;
@@ -51,6 +52,7 @@ class ArrowFlightDataStreamAuthenticationTest {
   private static final String SCHEMA_NAME = "secured-schema";
   private static final String USERNAME = "hop";
   private static final String PASSWORD = "s3cr3t";
+  private static final String LOOPBACK = "127.0.0.1";
 
   private final IVariables variables = new Variables();
   private final ILogChannel log = mock(ILogChannel.class);
@@ -74,8 +76,7 @@ class ArrowFlightDataStreamAuthenticationTest {
       reader.close();
     }
     if (server != null) {
-      server.shutdown();
-      server.getFlightServer().awaitTermination();
+      server.close();
     }
   }
 
@@ -97,7 +98,7 @@ class ArrowFlightDataStreamAuthenticationTest {
 
     ArrowFlightDataStream dataStream = new ArrowFlightDataStream();
     dataStream.setSchemaDefinitionName(SCHEMA_NAME);
-    dataStream.setHostname("localhost");
+    dataStream.setHostname(LOOPBACK);
     dataStream.setUsername(USERNAME);
     dataStream.setPassword(password);
 
@@ -110,7 +111,7 @@ class ArrowFlightDataStreamAuthenticationTest {
     security.setUsername(USERNAME);
     security.setPassword(PASSWORD);
 
-    server = new ArrowFlightServer("localhost", 0, security, variables, 
metadataProvider, log);
+    server = new ArrowFlightServer(LOOPBACK, 0, security, variables, 
metadataProvider, log);
     server.start();
     dataStream.setPort(Integer.toString(server.getFlightServer().getPort()));
 
@@ -129,8 +130,10 @@ class ArrowFlightDataStreamAuthenticationTest {
         
metadataProvider.getSerializer(SchemaDefinition.class).load(SCHEMA_NAME).getRowMeta();
 
     // Write a handful of rows. This authenticates and then does a doPut.
+    // The metadata object stays configuration: the writer is a copy, so the 
server can read
+    // that configuration without sharing the writer's allocator, buffer, or 
client.
     //
-    writer = (ArrowFlightDataStream) dataStreamMeta.getDataStream();
+    writer = ((ArrowFlightDataStream) dataStreamMeta.getDataStream()).clone();
     writer.initialize(variables, metadataProvider, true, dataStreamMeta);
     writer.setRowMeta(rowMeta);
     for (int i = 0; i < 5; i++) {
@@ -141,8 +144,7 @@ class ArrowFlightDataStreamAuthenticationTest {
     // Read them back. This authenticates again and then does a getFlightInfo 
and a doGet.
     //
     DataStreamMeta readMeta = loadStreamMeta();
-    reader = (ArrowFlightDataStream) readMeta.getDataStream();
-    reader.setPort(Integer.toString(server.getFlightServer().getPort()));
+    reader = ((ArrowFlightDataStream) readMeta.getDataStream()).clone();
     reader.initialize(variables, metadataProvider, false, readMeta);
 
     IRowMeta readRowMeta = reader.getRowMeta();
@@ -163,9 +165,11 @@ class ArrowFlightDataStreamAuthenticationTest {
     IRowMeta rowMeta =
         
metadataProvider.getSerializer(SchemaDefinition.class).load(SCHEMA_NAME).getRowMeta();
 
-    writer = (ArrowFlightDataStream) dataStreamMeta.getDataStream();
+    writer = ((ArrowFlightDataStream) dataStreamMeta.getDataStream()).clone();
     writer.initialize(variables, metadataProvider, true, dataStreamMeta);
 
-    assertThrows(HopException.class, () -> writer.setRowMeta(rowMeta));
+    ExpectedFlightRejection.run(
+        "data stream password does not match the Flight server",
+        () -> assertThrows(HopException.class, () -> 
writer.setRowMeta(rowMeta)));
   }
 }
diff --git 
a/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/flight/ArrowFlightServerAuthenticationTest.java
 
b/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/flight/ArrowFlightServerAuthenticationTest.java
index d9ca935811..86ec13539e 100644
--- 
a/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/flight/ArrowFlightServerAuthenticationTest.java
+++ 
b/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/flight/ArrowFlightServerAuthenticationTest.java
@@ -48,6 +48,7 @@ class ArrowFlightServerAuthenticationTest {
 
   private static final String USERNAME = "hop";
   private static final String PASSWORD = "s3cr3t";
+  private static final String LOOPBACK = "127.0.0.1";
 
   private final ILogChannel log = mock(ILogChannel.class);
 
@@ -64,8 +65,7 @@ class ArrowFlightServerAuthenticationTest {
       clientAllocator.close();
     }
     if (server != null) {
-      server.shutdown();
-      server.getFlightServer().awaitTermination();
+      server.close();
     }
   }
 
@@ -74,14 +74,12 @@ class ArrowFlightServerAuthenticationTest {
     //
     server =
         new ArrowFlightServer(
-            "localhost", 0, security, new Variables(), new 
MemoryMetadataProvider(), log);
+            LOOPBACK, 0, security, new Variables(), new 
MemoryMetadataProvider(), log);
     server.start();
 
     clientAllocator = new RootAllocator();
     client =
-        FlightClient.builder(
-                clientAllocator,
-                Location.forGrpcInsecure("localhost", 
server.getFlightServer().getPort()))
+        FlightClient.builder(clientAllocator, 
Location.forGrpcInsecure(LOOPBACK, server.getPort()))
             .build();
     return client;
   }
@@ -125,22 +123,30 @@ class ArrowFlightServerAuthenticationTest {
   void aServerWithCredentialsRejectsTheWrongPassword() throws Exception {
     FlightClient flightClient = startServerAndConnect(withCredentials());
 
-    FlightRuntimeException exception =
-        assertThrows(
-            FlightRuntimeException.class,
-            () -> flightClient.authenticateBasicToken(USERNAME, "wrong"));
-    assertEquals(FlightStatusCode.UNAUTHENTICATED, exception.status().code());
+    ExpectedFlightRejection.run(
+        "wrong password for user '" + USERNAME + "'",
+        () -> {
+          FlightRuntimeException exception =
+              assertThrows(
+                  FlightRuntimeException.class,
+                  () -> flightClient.authenticateBasicToken(USERNAME, 
"wrong"));
+          assertEquals(FlightStatusCode.UNAUTHENTICATED, 
exception.status().code());
+        });
   }
 
   @Test
   void aServerWithCredentialsRejectsTheWrongUsername() throws Exception {
     FlightClient flightClient = startServerAndConnect(withCredentials());
 
-    FlightRuntimeException exception =
-        assertThrows(
-            FlightRuntimeException.class,
-            () -> flightClient.authenticateBasicToken("somebody-else", 
PASSWORD));
-    assertEquals(FlightStatusCode.UNAUTHENTICATED, exception.status().code());
+    ExpectedFlightRejection.run(
+        "unknown user 'somebody-else'",
+        () -> {
+          FlightRuntimeException exception =
+              assertThrows(
+                  FlightRuntimeException.class,
+                  () -> flightClient.authenticateBasicToken("somebody-else", 
PASSWORD));
+          assertEquals(FlightStatusCode.UNAUTHENTICATED, 
exception.status().code());
+        });
   }
 
   @Test
diff --git 
a/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/flight/ExpectedFlightRejection.java
 
b/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/flight/ExpectedFlightRejection.java
new file mode 100644
index 0000000000..2435f15b25
--- /dev/null
+++ 
b/plugins/tech/arrow/src/test/java/org/apache/hop/arrow/flight/ExpectedFlightRejection.java
@@ -0,0 +1,56 @@
+/*
+ * 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.hop.arrow.flight;
+
+import org.apache.logging.log4j.Level;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.core.Logger;
+
+/**
+ * Runs a call that is expected to be rejected by the Flight server. Arrow's 
handshake logs that
+ * rejection at ERROR, with a stack trace, before the exception reaches the 
test. The log line below
+ * is the signal that this was intentional.
+ */
+public final class ExpectedFlightRejection {
+
+  private static final String HANDSHAKE_LOGGER =
+      "org.apache.arrow.flight.auth2.ClientHandshakeWrapper";
+
+  private ExpectedFlightRejection() {}
+
+  public static void run(String reason, RejectionCall call) throws Exception {
+    System.out.println(
+        "[test] Expected UNAUTHENTICATED from the Arrow Flight server: "
+            + reason
+            + ". Arrow would log this rejection at ERROR; that stack trace is 
hidden.");
+    Logger handshakeLogger = (Logger) LogManager.getLogger(HANDSHAKE_LOGGER);
+    Level previousLevel = handshakeLogger.getLevel();
+    handshakeLogger.setLevel(Level.OFF);
+    try {
+      call.run();
+    } finally {
+      handshakeLogger.setLevel(previousLevel);
+    }
+  }
+
+  @FunctionalInterface
+  public interface RejectionCall {
+    void run() throws Exception;
+  }
+}

Reply via email to