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;
+ }
+}