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

lidavidm pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-adbc.git


The following commit(s) were added to refs/heads/main by this push:
     new a3a94d4b7 fix(java/driver/flight-sql): include connectionOptions for 
preparedStatement.close() (#4513)
a3a94d4b7 is described below

commit a3a94d4b79c145f493801e8517d808a8d8a6ebd4
Author: Daniel_McBride <[email protected]>
AuthorDate: Sat Jul 18 18:24:25 2026 -0700

    fix(java/driver/flight-sql): include connectionOptions for 
preparedStatement.close() (#4513)
    
    Proposed solution for #4512
    
    Include connection options when calling preparedStatement.close()
    
    ---------
    
    Co-authored-by: dmcbride <[email protected]>
---
 .../adbc/driver/flightsql/FlightSqlConnection.java |  4 +--
 .../adbc/driver/flightsql/FlightSqlStatement.java  | 16 +++++----
 .../arrow/adbc/driver/flightsql/HeaderTest.java    | 39 +++++++++++++++++++++-
 .../adbc/driver/flightsql/HeaderValidator.java     |  4 +++
 4 files changed, 54 insertions(+), 9 deletions(-)

diff --git 
a/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java
 
b/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java
index 123b9828b..99c9fb735 100644
--- 
a/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java
+++ 
b/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java
@@ -126,7 +126,7 @@ public class FlightSqlConnection implements AdbcConnection {
 
   @Override
   public AdbcStatement createStatement() throws AdbcException {
-    return new FlightSqlStatement(allocator, client, clientCache, quirks);
+    return new FlightSqlStatement(allocator, client, clientCache, quirks, 
callOptions);
   }
 
   @Override
@@ -157,7 +157,7 @@ public class FlightSqlConnection implements AdbcConnection {
   public AdbcStatement bulkIngest(String targetTableName, BulkIngestMode mode)
       throws AdbcException {
     return FlightSqlStatement.ingestRoot(
-        allocator, client, clientCache, quirks, targetTableName, mode);
+        allocator, client, clientCache, quirks, targetTableName, mode, 
callOptions);
   }
 
   @Override
diff --git 
a/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlStatement.java
 
b/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlStatement.java
index dd3ba1f28..557353c74 100644
--- 
a/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlStatement.java
+++ 
b/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlStatement.java
@@ -29,6 +29,7 @@ import org.apache.arrow.adbc.core.AdbcStatusCode;
 import org.apache.arrow.adbc.core.BulkIngestMode;
 import org.apache.arrow.adbc.core.PartitionDescriptor;
 import org.apache.arrow.adbc.sql.SqlQuirks;
+import org.apache.arrow.flight.CallOption;
 import org.apache.arrow.flight.FlightEndpoint;
 import org.apache.arrow.flight.FlightInfo;
 import org.apache.arrow.flight.FlightRuntimeException;
@@ -36,7 +37,6 @@ import org.apache.arrow.flight.Location;
 import org.apache.arrow.flight.impl.Flight;
 import org.apache.arrow.flight.sql.FlightSqlClient;
 import org.apache.arrow.memory.BufferAllocator;
-import org.apache.arrow.util.AutoCloseables;
 import org.apache.arrow.vector.VectorSchemaRoot;
 import org.apache.arrow.vector.types.pojo.Field;
 import org.apache.arrow.vector.types.pojo.Schema;
@@ -47,6 +47,7 @@ public class FlightSqlStatement implements AdbcStatement {
   private final FlightSqlClientWithCallOptions client;
   private final LoadingCache<Location, FlightSqlClientWithCallOptions> 
clientCache;
   private final SqlQuirks quirks;
+  private final CallOption[] connectionOptions;
 
   // State for SQL queries
   private @Nullable String sqlQuery;
@@ -59,7 +60,8 @@ public class FlightSqlStatement implements AdbcStatement {
       BufferAllocator allocator,
       FlightSqlClientWithCallOptions client,
       LoadingCache<Location, FlightSqlClientWithCallOptions> clientCache,
-      SqlQuirks quirks) {
+      SqlQuirks quirks,
+      CallOption... connectionOptions) {
     this.allocator = allocator;
     this.client = client;
     this.clientCache = clientCache;
@@ -68,6 +70,7 @@ public class FlightSqlStatement implements AdbcStatement {
     this.preparedStatement = null;
     this.bulkOperation = null;
     this.bindRoot = null;
+    this.connectionOptions = connectionOptions;
   }
 
   static FlightSqlStatement ingestRoot(
@@ -76,10 +79,11 @@ public class FlightSqlStatement implements AdbcStatement {
       LoadingCache<Location, FlightSqlClientWithCallOptions> clientCache,
       SqlQuirks quirks,
       String targetTableName,
-      BulkIngestMode mode) {
+      BulkIngestMode mode,
+      CallOption... connectionOptions) {
     Objects.requireNonNull(targetTableName);
     final FlightSqlStatement statement =
-        new FlightSqlStatement(allocator, client, clientCache, quirks);
+        new FlightSqlStatement(allocator, client, clientCache, quirks, 
connectionOptions);
     statement.bulkOperation = new BulkState(mode, targetTableName);
     return statement;
   }
@@ -170,7 +174,7 @@ public class FlightSqlStatement implements AdbcStatement {
         statement.setParameters(new NonOwningRoot(bindParams));
         client.executePreparedUpdate(statement);
       } finally {
-        statement.close();
+        statement.close(connectionOptions);
       }
     } catch (FlightRuntimeException e) {
       // XXX: FlightSqlClient.executeUpdate does some extra wrapping that we 
need to undo
@@ -313,7 +317,7 @@ public class FlightSqlStatement implements AdbcStatement {
     // TODO(https://github.com/apache/arrow/issues/39814): this is annotated 
wrongly upstream
     if (preparedStatement != null) {
       try {
-        AutoCloseables.close(preparedStatement);
+        preparedStatement.close(connectionOptions);
       } catch (Exception e) {
         throw AdbcException.internal("[Flight SQL] Could not close prepared 
statement")
             .withCause(e);
diff --git 
a/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderTest.java
 
b/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderTest.java
index 261d172e8..b0a25758a 100644
--- 
a/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderTest.java
+++ 
b/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderTest.java
@@ -27,6 +27,7 @@ import org.apache.arrow.adbc.core.AdbcDatabase;
 import org.apache.arrow.adbc.core.AdbcDriver;
 import org.apache.arrow.adbc.core.AdbcException;
 import org.apache.arrow.adbc.core.AdbcInfoCode;
+import org.apache.arrow.adbc.core.AdbcStatement;
 import org.apache.arrow.adbc.core.AdbcStatusCode;
 import org.apache.arrow.adbc.drivermanager.AdbcDriverManager;
 import org.apache.arrow.driver.jdbc.utils.MockFlightSqlProducer;
@@ -55,16 +56,18 @@ public class HeaderTest {
   private AdbcConnection connection;
   private BufferAllocator allocator;
   private HeaderValidator.Factory headerValidatorFactory;
+  private MockFlightSqlProducer producer;
 
   @BeforeEach
   public void setUp() {
     allocator = new RootAllocator(Long.MAX_VALUE);
     headerValidatorFactory = new HeaderValidator.Factory();
+    producer = new MockFlightSqlProducer();
     builder =
         FlightServer.builder()
             .middleware(HeaderValidator.KEY, headerValidatorFactory)
             .location(Location.forGrpcInsecure("localhost", 0))
-            .producer(new MockFlightSqlProducer());
+            .producer(producer);
     params = new HashMap<>();
   }
 
@@ -89,6 +92,40 @@ public class HeaderTest {
     assertEquals(dummyValue, headers.get(dummyHeaderName));
   }
 
+  /**
+   * The connection's call options (which carry arbitrary headers, auth, 
cookies, etc.) must ride on
+   * every RPC the driver issues on that connection's behalf, including the 
ClosePreparedStatement
+   * call made when a prepared statement is closed. This regression test 
guards against dropping the
+   * connection options on close.
+   */
+  @Test
+  public void testHeaderSentWhenClosingPreparedStatement() throws Exception {
+    final String dummyValue = "dummy";
+    final String dummyHeaderName = "test-header";
+    final String query = "UPDATE the_table SET x = 1";
+    params.put(FlightSqlConnectionProperties.RPC_CALL_HEADER_PREFIX + 
dummyHeaderName, dummyValue);
+    producer.addUpdateQuery(query, /* updatedRows= */ 1L);
+    server = builder.build();
+    server.start();
+    connect();
+
+    // Prepare and then close a statement. Closing triggers a 
ClosePreparedStatement RPC, which must
+    // still carry the connection header. Before the fix, close() dropped the 
call options.
+    try (AdbcStatement statement = connection.createStatement()) {
+      statement.setSqlQuery(query);
+      statement.prepare();
+    }
+
+    // The connection header must appear on every RPC, including the final 
ClosePreparedStatement.
+    final int requestCount = headerValidatorFactory.getRequestCount();
+    assertTrue(requestCount > 0);
+    for (int i = 0; i < requestCount; i++) {
+      CallHeaders headers = 
headerValidatorFactory.getHeadersReceivedAtRequest(i);
+      assertEquals(
+          dummyValue, headers.get(dummyHeaderName), "connection header missing 
on RPC #" + i);
+    }
+  }
+
   @Test
   public void testCookies() throws Exception {
     builder.middleware(CookieMiddleware.KEY, new CookieMiddleware.Factory());
diff --git 
a/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderValidator.java
 
b/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderValidator.java
index f543dd995..c4702eeed 100644
--- 
a/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderValidator.java
+++ 
b/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderValidator.java
@@ -52,6 +52,10 @@ public class HeaderValidator implements 
FlightServerMiddleware {
       return cloneHeaders(headersReceived.get(request));
     }
 
+    public int getRequestCount() {
+      return headersReceived.size();
+    }
+
     private static CallHeaders cloneHeaders(CallHeaders headers) {
       FlightCallHeaders cloneHeaders = new FlightCallHeaders();
       for (String key : headers.keys()) {

Reply via email to