Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -262,15 +262,83 @@ public FlightInfo getInfo(final String query) {
@Override
public void close() throws SQLException {
if (catalog.isPresent()) {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
try {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to close Flight SQL session.", "closing Flight SQL session");
}
}
try {
AutoCloseables.close(sqlClient);
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to clean up client resources.", "closing Flight SQL client");
} catch (final Exception e) {
throw new SQLException("Failed to clean up client resources.", e);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures.
*
* @param fre the FlightRuntimeException to handle
* @param sqlErrorMessage the SQLException message to use for genuine failures
* @param operationDescription description of the operation for logging
* @throws SQLException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String sqlErrorMessage, String operationDescription) throws SQLException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw new SQLException(sqlErrorMessage, fre);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures as FlightRuntimeException.
*
* @param fre the FlightRuntimeException to handle
* @param operationDescription description of the operation for logging
* @throws FlightRuntimeException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String operationDescription) throws FlightRuntimeException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw fre;
}
}

/**
* Determines if a FlightRuntimeException represents a benign close operation error
* that should be suppressed.
*
* @param fre the FlightRuntimeException to check
* @return true if the exception should be suppressed, false otherwise
*/
private boolean isBenignCloseException(FlightRuntimeException fre) {
return fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage() != null
&& fre.getMessage().contains("Connection closed after GOAWAY"));
}

/**
* Logs a suppressed close exception with appropriate level based on debug settings.
*
* @param fre the FlightRuntimeException being suppressed
* @param operationDescription description of the operation for logging
*/
private void logSuppressedCloseException(FlightRuntimeException fre, String operationDescription) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer during shutdown
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Suppressed error {}", operationDescription, fre);
} else {
LOGGER.info("Suppressed benign error {}: {}", operationDescription, fre.getMessage());
}
Comment on lines +334 to +339

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why the switch? Just always put it at debug level.

}

/** A prepared statement handler. */
public interface PreparedStatement extends AutoCloseable {
/**
Expand DownExpand Up@@ -386,14 +454,7 @@ public void close() {
try {
preparedStatement.close(getOptions());
} catch (FlightRuntimeException fre) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer
if (fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage().contains("Connection closed after GOAWAY"))) {
LOGGER.warn("Supressed error closing PreparedStatement", fre);
return;
}
throw fre;
handleBenignCloseException(fre, "closing PreparedStatement");
}
}
};
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
/*
* 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.arrow.driver.jdbc.client;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertTrue;

import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.apache.arrow.driver.jdbc.FlightServerTestExtension;
import org.apache.arrow.flight.sql.NoOpFlightSqlProducer;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.util.AutoCloseables;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.slf4j.LoggerFactory;

/** Integration tests for {@link ArrowFlightSqlClientHandler} error suppression functionality. */
public class ArrowFlightSqlClientHandlerIntegrationTest {

/** Minimal producer for integration tests. */
public static class TestFlightSqlProducer extends NoOpFlightSqlProducer {
// No custom behavior needed for these tests
}

private static final TestFlightSqlProducer PRODUCER = new TestFlightSqlProducer();

@RegisterExtension
public static final FlightServerTestExtension FLIGHT_SERVER_TEST_EXTENSION =
FlightServerTestExtension.createStandardTestExtension(PRODUCER);

private static BufferAllocator allocator;
private Logger logger;
private ListAppender<ILoggingEvent> logAppender;

@BeforeAll
public static void setup() {
allocator = new RootAllocator(Long.MAX_VALUE);
}

@AfterAll
public static void tearDown() throws Exception {
AutoCloseables.close(PRODUCER, allocator);
}

@BeforeEach
public void setUp() {
// Set up logging capture
logger = (Logger) LoggerFactory.getLogger(ArrowFlightSqlClientHandler.class);
logAppender = new ListAppender<>();
logAppender.start();
logger.addAppender(logAppender);
logger.setLevel(Level.DEBUG);
}

@AfterEach
public void tearDownLogging() {
if (logger != null && logAppender != null) {
logger.detachAppender(logAppender);
}
}

// Note: Integration tests for closeSession() with catalog are not included because
// closeSession is a gRPC service method that's not routed through the FlightProducer,
// making it difficult to simulate errors in a test environment. The unit tests
// (ArrowFlightSqlClientHandlerTest) provide comprehensive coverage of the error
// suppression logic using reflection to test the private methods directly.

@Test
public void testClose_WithoutCatalog_NoCloseSessionCall() throws Exception {
// Arrange - no catalog means no CloseSession RPC
try (ArrowFlightSqlClientHandler client = new ArrowFlightSqlClientHandler.Builder()
.withHost(FLIGHT_SERVER_TEST_EXTENSION.getHost())
.withPort(FLIGHT_SERVER_TEST_EXTENSION.getPort())
.withBufferAllocator(allocator)
.withEncryption(false)
// No catalog set
.build()) {

// Act & Assert - should close successfully without any CloseSession RPC
assertDoesNotThrow(() -> client.close());
}

// Verify no CloseSession-related logging occurred
assertTrue(logAppender.list.stream()
.noneMatch(event -> event.getFormattedMessage().contains("closing Flight SQL session")));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -262,15 +262,83 @@ public FlightInfo getInfo(final String query) {
@Override
public void close() throws SQLException {
if (catalog.isPresent()) {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
try {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to close Flight SQL session.", "closing Flight SQL session");
}
}
try {
AutoCloseables.close(sqlClient);
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to clean up client resources.", "closing Flight SQL client");
} catch (final Exception e) {
throw new SQLException("Failed to clean up client resources.", e);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures.
*
* @param fre the FlightRuntimeException to handle
* @param sqlErrorMessage the SQLException message to use for genuine failures
* @param operationDescription description of the operation for logging
* @throws SQLException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String sqlErrorMessage, String operationDescription) throws SQLException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw new SQLException(sqlErrorMessage, fre);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures as FlightRuntimeException.
*
* @param fre the FlightRuntimeException to handle
* @param operationDescription description of the operation for logging
* @throws FlightRuntimeException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String operationDescription) throws FlightRuntimeException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw fre;
}
}

/**
* Determines if a FlightRuntimeException represents a benign close operation error
* that should be suppressed.
*
* @param fre the FlightRuntimeException to check
* @return true if the exception should be suppressed, false otherwise
*/
private boolean isBenignCloseException(FlightRuntimeException fre) {
return fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage() != null
&& fre.getMessage().contains("Connection closed after GOAWAY"));
}

/**
* Logs a suppressed close exception with appropriate level based on debug settings.
*
* @param fre the FlightRuntimeException being suppressed
* @param operationDescription description of the operation for logging
*/
private void logSuppressedCloseException(FlightRuntimeException fre, String operationDescription) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer during shutdown
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Suppressed error {}", operationDescription, fre);
} else {
LOGGER.info("Suppressed benign error {}: {}", operationDescription, fre.getMessage());
}
Comment on lines +334 to +339

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why the switch? Just always put it at debug level.

}

/** A prepared statement handler. */
public interface PreparedStatement extends AutoCloseable {
/**
Expand DownExpand Up@@ -386,14 +454,7 @@ public void close() {
try {
preparedStatement.close(getOptions());
} catch (FlightRuntimeException fre) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer
if (fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage().contains("Connection closed after GOAWAY"))) {
LOGGER.warn("Supressed error closing PreparedStatement", fre);
return;
}
throw fre;
handleBenignCloseException(fre, "closing PreparedStatement");
}
}
};
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
/*
* 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.arrow.driver.jdbc.client;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertTrue;

import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.apache.arrow.driver.jdbc.FlightServerTestExtension;
import org.apache.arrow.flight.sql.NoOpFlightSqlProducer;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.util.AutoCloseables;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.slf4j.LoggerFactory;

/** Integration tests for {@link ArrowFlightSqlClientHandler} error suppression functionality. */
public class ArrowFlightSqlClientHandlerIntegrationTest {

/** Minimal producer for integration tests. */
public static class TestFlightSqlProducer extends NoOpFlightSqlProducer {
// No custom behavior needed for these tests
}

private static final TestFlightSqlProducer PRODUCER = new TestFlightSqlProducer();

@RegisterExtension
public static final FlightServerTestExtension FLIGHT_SERVER_TEST_EXTENSION =
FlightServerTestExtension.createStandardTestExtension(PRODUCER);

private static BufferAllocator allocator;
private Logger logger;
private ListAppender<ILoggingEvent> logAppender;

@BeforeAll
public static void setup() {
allocator = new RootAllocator(Long.MAX_VALUE);
}

@AfterAll
public static void tearDown() throws Exception {
AutoCloseables.close(PRODUCER, allocator);
}

@BeforeEach
public void setUp() {
// Set up logging capture
logger = (Logger) LoggerFactory.getLogger(ArrowFlightSqlClientHandler.class);
logAppender = new ListAppender<>();
logAppender.start();
logger.addAppender(logAppender);
logger.setLevel(Level.DEBUG);
}

@AfterEach
public void tearDownLogging() {
if (logger != null && logAppender != null) {
logger.detachAppender(logAppender);
}
}

// Note: Integration tests for closeSession() with catalog are not included because
// closeSession is a gRPC service method that's not routed through the FlightProducer,
// making it difficult to simulate errors in a test environment. The unit tests
// (ArrowFlightSqlClientHandlerTest) provide comprehensive coverage of the error
// suppression logic using reflection to test the private methods directly.

@Test
public void testClose_WithoutCatalog_NoCloseSessionCall() throws Exception {
// Arrange - no catalog means no CloseSession RPC
try (ArrowFlightSqlClientHandler client = new ArrowFlightSqlClientHandler.Builder()
.withHost(FLIGHT_SERVER_TEST_EXTENSION.getHost())
.withPort(FLIGHT_SERVER_TEST_EXTENSION.getPort())
.withBufferAllocator(allocator)
.withEncryption(false)
// No catalog set
.build()) {

// Act & Assert - should close successfully without any CloseSession RPC
assertDoesNotThrow(() -> client.close());
}

// Verify no CloseSession-related logging occurred
assertTrue(logAppender.list.stream()
.noneMatch(event -> event.getFormattedMessage().contains("closing Flight SQL session")));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -262,15 +262,83 @@ public FlightInfo getInfo(final String query) {
@Override
public void close() throws SQLException {
if (catalog.isPresent()) {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
try {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to close Flight SQL session.", "closing Flight SQL session");
}
}
try {
AutoCloseables.close(sqlClient);
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to clean up client resources.", "closing Flight SQL client");
} catch (final Exception e) {
throw new SQLException("Failed to clean up client resources.", e);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures.
*
* @param fre the FlightRuntimeException to handle
* @param sqlErrorMessage the SQLException message to use for genuine failures
* @param operationDescription description of the operation for logging
* @throws SQLException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String sqlErrorMessage, String operationDescription) throws SQLException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw new SQLException(sqlErrorMessage, fre);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures as FlightRuntimeException.
*
* @param fre the FlightRuntimeException to handle
* @param operationDescription description of the operation for logging
* @throws FlightRuntimeException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String operationDescription) throws FlightRuntimeException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw fre;
}
}

/**
* Determines if a FlightRuntimeException represents a benign close operation error
* that should be suppressed.
*
* @param fre the FlightRuntimeException to check
* @return true if the exception should be suppressed, false otherwise
*/
private boolean isBenignCloseException(FlightRuntimeException fre) {
return fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage() != null
&& fre.getMessage().contains("Connection closed after GOAWAY"));
}

/**
* Logs a suppressed close exception with appropriate level based on debug settings.
*
* @param fre the FlightRuntimeException being suppressed
* @param operationDescription description of the operation for logging
*/
private void logSuppressedCloseException(FlightRuntimeException fre, String operationDescription) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer during shutdown
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Suppressed error {}", operationDescription, fre);
} else {
LOGGER.info("Suppressed benign error {}: {}", operationDescription, fre.getMessage());
}
Comment on lines +334 to +339

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why the switch? Just always put it at debug level.

}

/** A prepared statement handler. */
public interface PreparedStatement extends AutoCloseable {
/**
Expand DownExpand Up@@ -386,14 +454,7 @@ public void close() {
try {
preparedStatement.close(getOptions());
} catch (FlightRuntimeException fre) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer
if (fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage().contains("Connection closed after GOAWAY"))) {
LOGGER.warn("Supressed error closing PreparedStatement", fre);
return;
}
throw fre;
handleBenignCloseException(fre, "closing PreparedStatement");
}
}
};
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
/*
* 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.arrow.driver.jdbc.client;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertTrue;

import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.apache.arrow.driver.jdbc.FlightServerTestExtension;
import org.apache.arrow.flight.sql.NoOpFlightSqlProducer;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.util.AutoCloseables;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.slf4j.LoggerFactory;

/** Integration tests for {@link ArrowFlightSqlClientHandler} error suppression functionality. */
public class ArrowFlightSqlClientHandlerIntegrationTest {

/** Minimal producer for integration tests. */
public static class TestFlightSqlProducer extends NoOpFlightSqlProducer {
// No custom behavior needed for these tests
}

private static final TestFlightSqlProducer PRODUCER = new TestFlightSqlProducer();

@RegisterExtension
public static final FlightServerTestExtension FLIGHT_SERVER_TEST_EXTENSION =
FlightServerTestExtension.createStandardTestExtension(PRODUCER);

private static BufferAllocator allocator;
private Logger logger;
private ListAppender<ILoggingEvent> logAppender;

@BeforeAll
public static void setup() {
allocator = new RootAllocator(Long.MAX_VALUE);
}

@AfterAll
public static void tearDown() throws Exception {
AutoCloseables.close(PRODUCER, allocator);
}

@BeforeEach
public void setUp() {
// Set up logging capture
logger = (Logger) LoggerFactory.getLogger(ArrowFlightSqlClientHandler.class);
logAppender = new ListAppender<>();
logAppender.start();
logger.addAppender(logAppender);
logger.setLevel(Level.DEBUG);
}

@AfterEach
public void tearDownLogging() {
if (logger != null && logAppender != null) {
logger.detachAppender(logAppender);
}
}

// Note: Integration tests for closeSession() with catalog are not included because
// closeSession is a gRPC service method that's not routed through the FlightProducer,
// making it difficult to simulate errors in a test environment. The unit tests
// (ArrowFlightSqlClientHandlerTest) provide comprehensive coverage of the error
// suppression logic using reflection to test the private methods directly.

@Test
public void testClose_WithoutCatalog_NoCloseSessionCall() throws Exception {
// Arrange - no catalog means no CloseSession RPC
try (ArrowFlightSqlClientHandler client = new ArrowFlightSqlClientHandler.Builder()
.withHost(FLIGHT_SERVER_TEST_EXTENSION.getHost())
.withPort(FLIGHT_SERVER_TEST_EXTENSION.getPort())
.withBufferAllocator(allocator)
.withEncryption(false)
// No catalog set
.build()) {

// Act & Assert - should close successfully without any CloseSession RPC
assertDoesNotThrow(() -> client.close());
}

// Verify no CloseSession-related logging occurred
assertTrue(logAppender.list.stream()
.noneMatch(event -> event.getFormattedMessage().contains("closing Flight SQL session")));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -262,15 +262,83 @@ public FlightInfo getInfo(final String query) {
@Override
public void close() throws SQLException {
if (catalog.isPresent()) {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
try {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to close Flight SQL session.", "closing Flight SQL session");
}
}
try {
AutoCloseables.close(sqlClient);
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to clean up client resources.", "closing Flight SQL client");
} catch (final Exception e) {
throw new SQLException("Failed to clean up client resources.", e);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures.
*
* @param fre the FlightRuntimeException to handle
* @param sqlErrorMessage the SQLException message to use for genuine failures
* @param operationDescription description of the operation for logging
* @throws SQLException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String sqlErrorMessage, String operationDescription) throws SQLException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw new SQLException(sqlErrorMessage, fre);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures as FlightRuntimeException.
*
* @param fre the FlightRuntimeException to handle
* @param operationDescription description of the operation for logging
* @throws FlightRuntimeException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String operationDescription) throws FlightRuntimeException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw fre;
}
}

/**
* Determines if a FlightRuntimeException represents a benign close operation error
* that should be suppressed.
*
* @param fre the FlightRuntimeException to check
* @return true if the exception should be suppressed, false otherwise
*/
private boolean isBenignCloseException(FlightRuntimeException fre) {
return fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage() != null
&& fre.getMessage().contains("Connection closed after GOAWAY"));
}

/**
* Logs a suppressed close exception with appropriate level based on debug settings.
*
* @param fre the FlightRuntimeException being suppressed
* @param operationDescription description of the operation for logging
*/
private void logSuppressedCloseException(FlightRuntimeException fre, String operationDescription) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer during shutdown
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Suppressed error {}", operationDescription, fre);
} else {
LOGGER.info("Suppressed benign error {}: {}", operationDescription, fre.getMessage());
}
Comment on lines +334 to +339

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why the switch? Just always put it at debug level.

}

/** A prepared statement handler. */
public interface PreparedStatement extends AutoCloseable {
/**
Expand DownExpand Up@@ -386,14 +454,7 @@ public void close() {
try {
preparedStatement.close(getOptions());
} catch (FlightRuntimeException fre) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer
if (fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage().contains("Connection closed after GOAWAY"))) {
LOGGER.warn("Supressed error closing PreparedStatement", fre);
return;
}
throw fre;
handleBenignCloseException(fre, "closing PreparedStatement");
}
}
};
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
/*
* 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.arrow.driver.jdbc.client;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertTrue;

import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.apache.arrow.driver.jdbc.FlightServerTestExtension;
import org.apache.arrow.flight.sql.NoOpFlightSqlProducer;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.util.AutoCloseables;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.slf4j.LoggerFactory;

/** Integration tests for {@link ArrowFlightSqlClientHandler} error suppression functionality. */
public class ArrowFlightSqlClientHandlerIntegrationTest {

/** Minimal producer for integration tests. */
public static class TestFlightSqlProducer extends NoOpFlightSqlProducer {
// No custom behavior needed for these tests
}

private static final TestFlightSqlProducer PRODUCER = new TestFlightSqlProducer();

@RegisterExtension
public static final FlightServerTestExtension FLIGHT_SERVER_TEST_EXTENSION =
FlightServerTestExtension.createStandardTestExtension(PRODUCER);

private static BufferAllocator allocator;
private Logger logger;
private ListAppender<ILoggingEvent> logAppender;

@BeforeAll
public static void setup() {
allocator = new RootAllocator(Long.MAX_VALUE);
}

@AfterAll
public static void tearDown() throws Exception {
AutoCloseables.close(PRODUCER, allocator);
}

@BeforeEach
public void setUp() {
// Set up logging capture
logger = (Logger) LoggerFactory.getLogger(ArrowFlightSqlClientHandler.class);
logAppender = new ListAppender<>();
logAppender.start();
logger.addAppender(logAppender);
logger.setLevel(Level.DEBUG);
}

@AfterEach
public void tearDownLogging() {
if (logger != null && logAppender != null) {
logger.detachAppender(logAppender);
}
}

// Note: Integration tests for closeSession() with catalog are not included because
// closeSession is a gRPC service method that's not routed through the FlightProducer,
// making it difficult to simulate errors in a test environment. The unit tests
// (ArrowFlightSqlClientHandlerTest) provide comprehensive coverage of the error
// suppression logic using reflection to test the private methods directly.

@Test
public void testClose_WithoutCatalog_NoCloseSessionCall() throws Exception {
// Arrange - no catalog means no CloseSession RPC
try (ArrowFlightSqlClientHandler client = new ArrowFlightSqlClientHandler.Builder()
.withHost(FLIGHT_SERVER_TEST_EXTENSION.getHost())
.withPort(FLIGHT_SERVER_TEST_EXTENSION.getPort())
.withBufferAllocator(allocator)
.withEncryption(false)
// No catalog set
.build()) {

// Act & Assert - should close successfully without any CloseSession RPC
assertDoesNotThrow(() -> client.close());
}

// Verify no CloseSession-related logging occurred
assertTrue(logAppender.list.stream()
.noneMatch(event -> event.getFormattedMessage().contains("closing Flight SQL session")));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -262,15 +262,83 @@ public FlightInfo getInfo(final String query) {
@Override
public void close() throws SQLException {
if (catalog.isPresent()) {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
try {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to close Flight SQL session.", "closing Flight SQL session");
}
}
try {
AutoCloseables.close(sqlClient);
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to clean up client resources.", "closing Flight SQL client");
} catch (final Exception e) {
throw new SQLException("Failed to clean up client resources.", e);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures.
*
* @param fre the FlightRuntimeException to handle
* @param sqlErrorMessage the SQLException message to use for genuine failures
* @param operationDescription description of the operation for logging
* @throws SQLException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String sqlErrorMessage, String operationDescription) throws SQLException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw new SQLException(sqlErrorMessage, fre);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures as FlightRuntimeException.
*
* @param fre the FlightRuntimeException to handle
* @param operationDescription description of the operation for logging
* @throws FlightRuntimeException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String operationDescription) throws FlightRuntimeException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw fre;
}
}

/**
* Determines if a FlightRuntimeException represents a benign close operation error
* that should be suppressed.
*
* @param fre the FlightRuntimeException to check
* @return true if the exception should be suppressed, false otherwise
*/
private boolean isBenignCloseException(FlightRuntimeException fre) {
return fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage() != null
&& fre.getMessage().contains("Connection closed after GOAWAY"));
}

/**
* Logs a suppressed close exception with appropriate level based on debug settings.
*
* @param fre the FlightRuntimeException being suppressed
* @param operationDescription description of the operation for logging
*/
private void logSuppressedCloseException(FlightRuntimeException fre, String operationDescription) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer during shutdown
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Suppressed error {}", operationDescription, fre);
} else {
LOGGER.info("Suppressed benign error {}: {}", operationDescription, fre.getMessage());
}
Comment on lines +334 to +339

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why the switch? Just always put it at debug level.

}

/** A prepared statement handler. */
public interface PreparedStatement extends AutoCloseable {
/**
Expand DownExpand Up@@ -386,14 +454,7 @@ public void close() {
try {
preparedStatement.close(getOptions());
} catch (FlightRuntimeException fre) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer
if (fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage().contains("Connection closed after GOAWAY"))) {
LOGGER.warn("Supressed error closing PreparedStatement", fre);
return;
}
throw fre;
handleBenignCloseException(fre, "closing PreparedStatement");
}
}
};
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
/*
* 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.arrow.driver.jdbc.client;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertTrue;

import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.apache.arrow.driver.jdbc.FlightServerTestExtension;
import org.apache.arrow.flight.sql.NoOpFlightSqlProducer;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.util.AutoCloseables;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.slf4j.LoggerFactory;

/** Integration tests for {@link ArrowFlightSqlClientHandler} error suppression functionality. */
public class ArrowFlightSqlClientHandlerIntegrationTest {

/** Minimal producer for integration tests. */
public static class TestFlightSqlProducer extends NoOpFlightSqlProducer {
// No custom behavior needed for these tests
}

private static final TestFlightSqlProducer PRODUCER = new TestFlightSqlProducer();

@RegisterExtension
public static final FlightServerTestExtension FLIGHT_SERVER_TEST_EXTENSION =
FlightServerTestExtension.createStandardTestExtension(PRODUCER);

private static BufferAllocator allocator;
private Logger logger;
private ListAppender<ILoggingEvent> logAppender;

@BeforeAll
public static void setup() {
allocator = new RootAllocator(Long.MAX_VALUE);
}

@AfterAll
public static void tearDown() throws Exception {
AutoCloseables.close(PRODUCER, allocator);
}

@BeforeEach
public void setUp() {
// Set up logging capture
logger = (Logger) LoggerFactory.getLogger(ArrowFlightSqlClientHandler.class);
logAppender = new ListAppender<>();
logAppender.start();
logger.addAppender(logAppender);
logger.setLevel(Level.DEBUG);
}

@AfterEach
public void tearDownLogging() {
if (logger != null && logAppender != null) {
logger.detachAppender(logAppender);
}
}

// Note: Integration tests for closeSession() with catalog are not included because
// closeSession is a gRPC service method that's not routed through the FlightProducer,
// making it difficult to simulate errors in a test environment. The unit tests
// (ArrowFlightSqlClientHandlerTest) provide comprehensive coverage of the error
// suppression logic using reflection to test the private methods directly.

@Test
public void testClose_WithoutCatalog_NoCloseSessionCall() throws Exception {
// Arrange - no catalog means no CloseSession RPC
try (ArrowFlightSqlClientHandler client = new ArrowFlightSqlClientHandler.Builder()
.withHost(FLIGHT_SERVER_TEST_EXTENSION.getHost())
.withPort(FLIGHT_SERVER_TEST_EXTENSION.getPort())
.withBufferAllocator(allocator)
.withEncryption(false)
// No catalog set
.build()) {

// Act & Assert - should close successfully without any CloseSession RPC
assertDoesNotThrow(() -> client.close());
}

// Verify no CloseSession-related logging occurred
assertTrue(logAppender.list.stream()
.noneMatch(event -> event.getFormattedMessage().contains("closing Flight SQL session")));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -262,15 +262,83 @@ public FlightInfo getInfo(final String query) {
@Override
public void close() throws SQLException {
if (catalog.isPresent()) {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
try {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to close Flight SQL session.", "closing Flight SQL session");
}
}
try {
AutoCloseables.close(sqlClient);
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to clean up client resources.", "closing Flight SQL client");
} catch (final Exception e) {
throw new SQLException("Failed to clean up client resources.", e);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures.
*
* @param fre the FlightRuntimeException to handle
* @param sqlErrorMessage the SQLException message to use for genuine failures
* @param operationDescription description of the operation for logging
* @throws SQLException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String sqlErrorMessage, String operationDescription) throws SQLException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw new SQLException(sqlErrorMessage, fre);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures as FlightRuntimeException.
*
* @param fre the FlightRuntimeException to handle
* @param operationDescription description of the operation for logging
* @throws FlightRuntimeException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String operationDescription) throws FlightRuntimeException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw fre;
}
}

/**
* Determines if a FlightRuntimeException represents a benign close operation error
* that should be suppressed.
*
* @param fre the FlightRuntimeException to check
* @return true if the exception should be suppressed, false otherwise
*/
private boolean isBenignCloseException(FlightRuntimeException fre) {
return fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage() != null
&& fre.getMessage().contains("Connection closed after GOAWAY"));
}

/**
* Logs a suppressed close exception with appropriate level based on debug settings.
*
* @param fre the FlightRuntimeException being suppressed
* @param operationDescription description of the operation for logging
*/
private void logSuppressedCloseException(FlightRuntimeException fre, String operationDescription) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer during shutdown
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Suppressed error {}", operationDescription, fre);
} else {
LOGGER.info("Suppressed benign error {}: {}", operationDescription, fre.getMessage());
}
Comment on lines +334 to +339

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why the switch? Just always put it at debug level.

}

/** A prepared statement handler. */
public interface PreparedStatement extends AutoCloseable {
/**
Expand DownExpand Up@@ -386,14 +454,7 @@ public void close() {
try {
preparedStatement.close(getOptions());
} catch (FlightRuntimeException fre) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer
if (fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage().contains("Connection closed after GOAWAY"))) {
LOGGER.warn("Supressed error closing PreparedStatement", fre);
return;
}
throw fre;
handleBenignCloseException(fre, "closing PreparedStatement");
}
}
};
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
/*
* 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.arrow.driver.jdbc.client;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertTrue;

import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.apache.arrow.driver.jdbc.FlightServerTestExtension;
import org.apache.arrow.flight.sql.NoOpFlightSqlProducer;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.util.AutoCloseables;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.slf4j.LoggerFactory;

/** Integration tests for {@link ArrowFlightSqlClientHandler} error suppression functionality. */
public class ArrowFlightSqlClientHandlerIntegrationTest {

/** Minimal producer for integration tests. */
public static class TestFlightSqlProducer extends NoOpFlightSqlProducer {
// No custom behavior needed for these tests
}

private static final TestFlightSqlProducer PRODUCER = new TestFlightSqlProducer();

@RegisterExtension
public static final FlightServerTestExtension FLIGHT_SERVER_TEST_EXTENSION =
FlightServerTestExtension.createStandardTestExtension(PRODUCER);

private static BufferAllocator allocator;
private Logger logger;
private ListAppender<ILoggingEvent> logAppender;

@BeforeAll
public static void setup() {
allocator = new RootAllocator(Long.MAX_VALUE);
}

@AfterAll
public static void tearDown() throws Exception {
AutoCloseables.close(PRODUCER, allocator);
}

@BeforeEach
public void setUp() {
// Set up logging capture
logger = (Logger) LoggerFactory.getLogger(ArrowFlightSqlClientHandler.class);
logAppender = new ListAppender<>();
logAppender.start();
logger.addAppender(logAppender);
logger.setLevel(Level.DEBUG);
}

@AfterEach
public void tearDownLogging() {
if (logger != null && logAppender != null) {
logger.detachAppender(logAppender);
}
}

// Note: Integration tests for closeSession() with catalog are not included because
// closeSession is a gRPC service method that's not routed through the FlightProducer,
// making it difficult to simulate errors in a test environment. The unit tests
// (ArrowFlightSqlClientHandlerTest) provide comprehensive coverage of the error
// suppression logic using reflection to test the private methods directly.

@Test
public void testClose_WithoutCatalog_NoCloseSessionCall() throws Exception {
// Arrange - no catalog means no CloseSession RPC
try (ArrowFlightSqlClientHandler client = new ArrowFlightSqlClientHandler.Builder()
.withHost(FLIGHT_SERVER_TEST_EXTENSION.getHost())
.withPort(FLIGHT_SERVER_TEST_EXTENSION.getPort())
.withBufferAllocator(allocator)
.withEncryption(false)
// No catalog set
.build()) {

// Act & Assert - should close successfully without any CloseSession RPC
assertDoesNotThrow(() -> client.close());
}

// Verify no CloseSession-related logging occurred
assertTrue(logAppender.list.stream()
.noneMatch(event -> event.getFormattedMessage().contains("closing Flight SQL session")));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -262,15 +262,83 @@ public FlightInfo getInfo(final String query) {
@Override
public void close() throws SQLException {
if (catalog.isPresent()) {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
try {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to close Flight SQL session.", "closing Flight SQL session");
}
}
try {
AutoCloseables.close(sqlClient);
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to clean up client resources.", "closing Flight SQL client");
} catch (final Exception e) {
throw new SQLException("Failed to clean up client resources.", e);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures.
*
* @param fre the FlightRuntimeException to handle
* @param sqlErrorMessage the SQLException message to use for genuine failures
* @param operationDescription description of the operation for logging
* @throws SQLException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String sqlErrorMessage, String operationDescription) throws SQLException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw new SQLException(sqlErrorMessage, fre);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures as FlightRuntimeException.
*
* @param fre the FlightRuntimeException to handle
* @param operationDescription description of the operation for logging
* @throws FlightRuntimeException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String operationDescription) throws FlightRuntimeException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw fre;
}
}

/**
* Determines if a FlightRuntimeException represents a benign close operation error
* that should be suppressed.
*
* @param fre the FlightRuntimeException to check
* @return true if the exception should be suppressed, false otherwise
*/
private boolean isBenignCloseException(FlightRuntimeException fre) {
return fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage() != null
&& fre.getMessage().contains("Connection closed after GOAWAY"));
}

/**
* Logs a suppressed close exception with appropriate level based on debug settings.
*
* @param fre the FlightRuntimeException being suppressed
* @param operationDescription description of the operation for logging
*/
private void logSuppressedCloseException(FlightRuntimeException fre, String operationDescription) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer during shutdown
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Suppressed error {}", operationDescription, fre);
} else {
LOGGER.info("Suppressed benign error {}: {}", operationDescription, fre.getMessage());
}
Comment on lines +334 to +339

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why the switch? Just always put it at debug level.

}

/** A prepared statement handler. */
public interface PreparedStatement extends AutoCloseable {
/**
Expand DownExpand Up@@ -386,14 +454,7 @@ public void close() {
try {
preparedStatement.close(getOptions());
} catch (FlightRuntimeException fre) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer
if (fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage().contains("Connection closed after GOAWAY"))) {
LOGGER.warn("Supressed error closing PreparedStatement", fre);
return;
}
throw fre;
handleBenignCloseException(fre, "closing PreparedStatement");
}
}
};
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
/*
* 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.arrow.driver.jdbc.client;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertTrue;

import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.apache.arrow.driver.jdbc.FlightServerTestExtension;
import org.apache.arrow.flight.sql.NoOpFlightSqlProducer;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.util.AutoCloseables;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.slf4j.LoggerFactory;

/** Integration tests for {@link ArrowFlightSqlClientHandler} error suppression functionality. */
public class ArrowFlightSqlClientHandlerIntegrationTest {

/** Minimal producer for integration tests. */
public static class TestFlightSqlProducer extends NoOpFlightSqlProducer {
// No custom behavior needed for these tests
}

private static final TestFlightSqlProducer PRODUCER = new TestFlightSqlProducer();

@RegisterExtension
public static final FlightServerTestExtension FLIGHT_SERVER_TEST_EXTENSION =
FlightServerTestExtension.createStandardTestExtension(PRODUCER);

private static BufferAllocator allocator;
private Logger logger;
private ListAppender<ILoggingEvent> logAppender;

@BeforeAll
public static void setup() {
allocator = new RootAllocator(Long.MAX_VALUE);
}

@AfterAll
public static void tearDown() throws Exception {
AutoCloseables.close(PRODUCER, allocator);
}

@BeforeEach
public void setUp() {
// Set up logging capture
logger = (Logger) LoggerFactory.getLogger(ArrowFlightSqlClientHandler.class);
logAppender = new ListAppender<>();
logAppender.start();
logger.addAppender(logAppender);
logger.setLevel(Level.DEBUG);
}

@AfterEach
public void tearDownLogging() {
if (logger != null && logAppender != null) {
logger.detachAppender(logAppender);
}
}

// Note: Integration tests for closeSession() with catalog are not included because
// closeSession is a gRPC service method that's not routed through the FlightProducer,
// making it difficult to simulate errors in a test environment. The unit tests
// (ArrowFlightSqlClientHandlerTest) provide comprehensive coverage of the error
// suppression logic using reflection to test the private methods directly.

@Test
public void testClose_WithoutCatalog_NoCloseSessionCall() throws Exception {
// Arrange - no catalog means no CloseSession RPC
try (ArrowFlightSqlClientHandler client = new ArrowFlightSqlClientHandler.Builder()
.withHost(FLIGHT_SERVER_TEST_EXTENSION.getHost())
.withPort(FLIGHT_SERVER_TEST_EXTENSION.getPort())
.withBufferAllocator(allocator)
.withEncryption(false)
// No catalog set
.build()) {

// Act & Assert - should close successfully without any CloseSession RPC
assertDoesNotThrow(() -> client.close());
}

// Verify no CloseSession-related logging occurred
assertTrue(logAppender.list.stream()
.noneMatch(event -> event.getFormattedMessage().contains("closing Flight SQL session")));
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -262,15 +262,83 @@ public FlightInfo getInfo(final String query) {
@Override
public void close() throws SQLException {
if (catalog.isPresent()) {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
try {
sqlClient.closeSession(new CloseSessionRequest(), getOptions());
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to close Flight SQL session.", "closing Flight SQL session");
}
}
try {
AutoCloseables.close(sqlClient);
} catch (FlightRuntimeException fre) {
handleBenignCloseException(fre, "Failed to clean up client resources.", "closing Flight SQL client");
} catch (final Exception e) {
throw new SQLException("Failed to clean up client resources.", e);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures.
*
* @param fre the FlightRuntimeException to handle
* @param sqlErrorMessage the SQLException message to use for genuine failures
* @param operationDescription description of the operation for logging
* @throws SQLException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String sqlErrorMessage, String operationDescription) throws SQLException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw new SQLException(sqlErrorMessage, fre);
}
}

/**
* Handles FlightRuntimeException during close operations, suppressing benign gRPC shutdown errors
* while re-throwing genuine failures as FlightRuntimeException.
*
* @param fre the FlightRuntimeException to handle
* @param operationDescription description of the operation for logging
* @throws FlightRuntimeException if the exception represents a genuine failure
*/
private void handleBenignCloseException(FlightRuntimeException fre, String operationDescription) throws FlightRuntimeException {
if (isBenignCloseException(fre)) {
logSuppressedCloseException(fre, operationDescription);
} else {
throw fre;
}
}

/**
* Determines if a FlightRuntimeException represents a benign close operation error
* that should be suppressed.
*
* @param fre the FlightRuntimeException to check
* @return true if the exception should be suppressed, false otherwise
*/
private boolean isBenignCloseException(FlightRuntimeException fre) {
return fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage() != null
&& fre.getMessage().contains("Connection closed after GOAWAY"));
}

/**
* Logs a suppressed close exception with appropriate level based on debug settings.
*
* @param fre the FlightRuntimeException being suppressed
* @param operationDescription description of the operation for logging
*/
private void logSuppressedCloseException(FlightRuntimeException fre, String operationDescription) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer during shutdown
if (LOGGER.isDebugEnabled()) {
LOGGER.debug("Suppressed error {}", operationDescription, fre);
} else {
LOGGER.info("Suppressed benign error {}: {}", operationDescription, fre.getMessage());
}
Comment on lines +334 to +339

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why the switch? Just always put it at debug level.

}

/** A prepared statement handler. */
public interface PreparedStatement extends AutoCloseable {
/**
Expand DownExpand Up@@ -386,14 +454,7 @@ public void close() {
try {
preparedStatement.close(getOptions());
} catch (FlightRuntimeException fre) {
// ARROW-17785: suppress exceptions caused by flaky gRPC layer
if (fre.status().code().equals(FlightStatusCode.UNAVAILABLE)
|| (fre.status().code().equals(FlightStatusCode.INTERNAL)
&& fre.getMessage().contains("Connection closed after GOAWAY"))) {
LOGGER.warn("Supressed error closing PreparedStatement", fre);
return;
}
throw fre;
handleBenignCloseException(fre, "closing PreparedStatement");
}
}
};
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
/*
* 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.arrow.driver.jdbc.client;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertTrue;

import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.apache.arrow.driver.jdbc.FlightServerTestExtension;
import org.apache.arrow.flight.sql.NoOpFlightSqlProducer;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.memory.RootAllocator;
import org.apache.arrow.util.AutoCloseables;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.slf4j.LoggerFactory;

/** Integration tests for {@link ArrowFlightSqlClientHandler} error suppression functionality. */
public class ArrowFlightSqlClientHandlerIntegrationTest {

/** Minimal producer for integration tests. */
public static class TestFlightSqlProducer extends NoOpFlightSqlProducer {
// No custom behavior needed for these tests
}

private static final TestFlightSqlProducer PRODUCER = new TestFlightSqlProducer();

@RegisterExtension
public static final FlightServerTestExtension FLIGHT_SERVER_TEST_EXTENSION =
FlightServerTestExtension.createStandardTestExtension(PRODUCER);

private static BufferAllocator allocator;
private Logger logger;
private ListAppender<ILoggingEvent> logAppender;

@BeforeAll
public static void setup() {
allocator = new RootAllocator(Long.MAX_VALUE);
}

@AfterAll
public static void tearDown() throws Exception {
AutoCloseables.close(PRODUCER, allocator);
}

@BeforeEach
public void setUp() {
// Set up logging capture
logger = (Logger) LoggerFactory.getLogger(ArrowFlightSqlClientHandler.class);
logAppender = new ListAppender<>();
logAppender.start();
logger.addAppender(logAppender);
logger.setLevel(Level.DEBUG);
}

@AfterEach
public void tearDownLogging() {
if (logger != null && logAppender != null) {
logger.detachAppender(logAppender);
}
}

// Note: Integration tests for closeSession() with catalog are not included because
// closeSession is a gRPC service method that's not routed through the FlightProducer,
// making it difficult to simulate errors in a test environment. The unit tests
// (ArrowFlightSqlClientHandlerTest) provide comprehensive coverage of the error
// suppression logic using reflection to test the private methods directly.

@Test
public void testClose_WithoutCatalog_NoCloseSessionCall() throws Exception {
// Arrange - no catalog means no CloseSession RPC
try (ArrowFlightSqlClientHandler client = new ArrowFlightSqlClientHandler.Builder()
.withHost(FLIGHT_SERVER_TEST_EXTENSION.getHost())
.withPort(FLIGHT_SERVER_TEST_EXTENSION.getPort())
.withBufferAllocator(allocator)
.withEncryption(false)
// No catalog set
.build()) {

// Act & Assert - should close successfully without any CloseSession RPC
assertDoesNotThrow(() -> client.close());
}

// Verify no CloseSession-related logging occurred
assertTrue(logAppender.list.stream()
.noneMatch(event -> event.getFormattedMessage().contains("closing Flight SQL session")));
}
}
Loading
Loading