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

abhishekrb19 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git


The following commit(s) were added to refs/heads/master by this push:
     new 2812445723d feat: Add remote address tracking for JDBC/Avatica SQL 
queries (#19231)
2812445723d is described below

commit 2812445723d4d11c407a2dfa790b973c7eb428cb
Author: ahrilee <[email protected]>
AuthorDate: Thu Jul 16 19:36:22 2026 +0900

    feat: Add remote address tracking for JDBC/Avatica SQL queries (#19231)
    
    fix: Collect remoteAddress for JDBC/Avatica query metrics
    
    sqlQuery/time, sqlQuery/bytes, and sqlQuery/planningTimeMs metrics now 
include remoteAddress for JDBC/Avatica queries, consistent with the HTTP SQL 
API.
    
    ---------
    
    Co-authored-by: egg.boiled <[email protected]>
---
 .../java/org/apache/druid/sql/DirectStatement.java | 11 +--
 .../org/apache/druid/sql/PreparedStatement.java    | 16 +++-
 .../org/apache/druid/sql/SqlStatementFactory.java  | 15 +++-
 .../druid/sql/avatica/DruidAvaticaJsonHandler.java |  6 ++
 .../sql/avatica/DruidAvaticaProtobufHandler.java   | 72 +++++++++--------
 .../apache/druid/sql/avatica/DruidConnection.java  |  6 +-
 .../sql/avatica/DruidJdbcPreparedStatement.java    |  5 +-
 .../druid/sql/avatica/DruidJdbcStatement.java      | 11 ++-
 .../org/apache/druid/sql/avatica/DruidMeta.java    | 25 +++++-
 .../org/apache/druid/sql/SqlStatementTest.java     |  6 +-
 .../druid/sql/avatica/DruidAvaticaHandlerTest.java | 89 ++++++++++++++++++++++
 .../sql/avatica/DruidAvaticaJsonHandlerTest.java   | 10 ++-
 .../avatica/DruidAvaticaProtobufHandlerTest.java   | 18 +++--
 .../druid/sql/avatica/DruidStatementTest.java      | 34 ++++-----
 .../sql/calcite/util/QueryFrameworkUtils.java      |  6 +-
 15 files changed, 242 insertions(+), 88 deletions(-)

diff --git a/sql/src/main/java/org/apache/druid/sql/DirectStatement.java 
b/sql/src/main/java/org/apache/druid/sql/DirectStatement.java
index 44f8351c5b2..f206194c732 100644
--- a/sql/src/main/java/org/apache/druid/sql/DirectStatement.java
+++ b/sql/src/main/java/org/apache/druid/sql/DirectStatement.java
@@ -34,6 +34,7 @@ import org.apache.druid.sql.calcite.planner.DruidPlanner;
 import org.apache.druid.sql.calcite.planner.PlannerResult;
 import org.apache.druid.sql.calcite.planner.PrepareResult;
 
+import javax.annotation.Nullable;
 import java.util.Set;
 
 /**
@@ -155,20 +156,12 @@ public class DirectStatement extends AbstractStatement 
implements Cancelable
   public DirectStatement(
       final SqlToolbox lifecycleToolbox,
       final SqlQueryPlus queryPlus,
-      final String remoteAddress
+      @Nullable final String remoteAddress
   )
   {
     super(lifecycleToolbox, queryPlus, remoteAddress);
   }
 
-  public DirectStatement(
-      final SqlToolbox lifecycleToolbox,
-      final SqlQueryPlus sqlRequest
-  )
-  {
-    super(lifecycleToolbox, sqlRequest, null);
-  }
-
   /**
    * Convenience method to perform Direct execution of a query. Does both
    * the {@link #plan()} step and the {@link ResultSet#run()} step.
diff --git a/sql/src/main/java/org/apache/druid/sql/PreparedStatement.java 
b/sql/src/main/java/org/apache/druid/sql/PreparedStatement.java
index 23737eee9a8..7fa6def5879 100644
--- a/sql/src/main/java/org/apache/druid/sql/PreparedStatement.java
+++ b/sql/src/main/java/org/apache/druid/sql/PreparedStatement.java
@@ -23,6 +23,7 @@ import org.apache.calcite.avatica.remote.TypedValue;
 import org.apache.druid.sql.calcite.planner.DruidPlanner;
 import org.apache.druid.sql.calcite.planner.PrepareResult;
 
+import javax.annotation.Nullable;
 import java.util.List;
 
 /**
@@ -34,10 +35,11 @@ public class PreparedStatement extends AbstractStatement
 
   public PreparedStatement(
       final SqlToolbox lifecycleToolbox,
-      final SqlQueryPlus queryPlus
+      final SqlQueryPlus queryPlus,
+      @Nullable final String remoteAddress
   )
   {
-    super(lifecycleToolbox, queryPlus, null);
+    super(lifecycleToolbox, queryPlus, remoteAddress);
     this.originalRequest = queryPlus;
   }
 
@@ -85,12 +87,18 @@ public class PreparedStatement extends AbstractStatement
    * same statement can be execute many times, including concurrently. Each
    * execution repeats the parse, validate, authorize and plan steps since
    * data, permissions, views and other dependencies may have changed.
+   * <p>
+   * The {@code remoteAddress} is supplied per execution (read from the execute
+   * request) rather than reusing the address captured at prepare time, so 
that a
+   * prepared statement reused after a reconnect or from a different remote is
+   * attributed to the actual caller of this execution.
    */
-  public DirectStatement execute(List<TypedValue> parameters)
+  public DirectStatement execute(List<TypedValue> parameters, @Nullable String 
remoteAddress)
   {
     return new DirectStatement(
         sqlToolbox,
-        originalRequest.freshCopy().withParameters(parameters)
+        originalRequest.freshCopy().withParameters(parameters),
+        remoteAddress
     );
   }
 
diff --git a/sql/src/main/java/org/apache/druid/sql/SqlStatementFactory.java 
b/sql/src/main/java/org/apache/druid/sql/SqlStatementFactory.java
index 9d31441befb..34fc2acfbd2 100644
--- a/sql/src/main/java/org/apache/druid/sql/SqlStatementFactory.java
+++ b/sql/src/main/java/org/apache/druid/sql/SqlStatementFactory.java
@@ -19,6 +19,7 @@
 
 package org.apache.druid.sql;
 
+import javax.annotation.Nullable;
 import javax.servlet.http.HttpServletRequest;
 
 /**
@@ -51,11 +52,21 @@ public class SqlStatementFactory
 
   public DirectStatement directStatement(final SqlQueryPlus sqlRequest)
   {
-    return new DirectStatement(lifecycleToolbox, sqlRequest);
+    return new DirectStatement(lifecycleToolbox, sqlRequest, null);
+  }
+
+  public DirectStatement directStatement(final SqlQueryPlus sqlRequest, 
@Nullable final String remoteAddress)
+  {
+    return new DirectStatement(lifecycleToolbox, sqlRequest, remoteAddress);
   }
 
   public PreparedStatement preparedStatement(final SqlQueryPlus sqlRequest)
   {
-    return new PreparedStatement(lifecycleToolbox, sqlRequest);
+    return new PreparedStatement(lifecycleToolbox, sqlRequest, null);
+  }
+
+  public PreparedStatement preparedStatement(final SqlQueryPlus sqlRequest, 
@Nullable final String remoteAddress)
+  {
+    return new PreparedStatement(lifecycleToolbox, sqlRequest, remoteAddress);
   }
 }
diff --git 
a/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandler.java 
b/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandler.java
index f6d0738e232..2e80bfff4e7 100644
--- 
a/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandler.java
+++ 
b/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandler.java
@@ -66,6 +66,9 @@ public class DruidAvaticaJsonHandler extends 
DruidAvaticaHandler
   public boolean handle(Request request, Response response, Callback callback) 
throws Exception
   {
     String requestURI = request.getHttpURI().getPath();
+    String remoteAddr = Request.getRemoteAddr(request);
+    DruidMeta.setThreadLocalRemoteAddress(remoteAddr);
+
     try (Timer.Context ctx = this.requestTimer.start()) {
       if 
(AVATICA_PATH_NO_TRAILING_SLASH.equals(StringUtils.maybeRemoveTrailingSlash(requestURI)))
 {
         response.getHeaders().put("Content-Type", 
"application/json;charset=utf-8");
@@ -114,6 +117,9 @@ public class DruidAvaticaJsonHandler extends 
DruidAvaticaHandler
         return true;
       }
     }
+    finally {
+      DruidMeta.clearThreadLocalRemoteAddress();
+    }
     return false;
   }
 
diff --git 
a/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandler.java
 
b/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandler.java
index 7d5c6183025..1aee0db7729 100644
--- 
a/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandler.java
+++ 
b/sql/src/main/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandler.java
@@ -68,43 +68,51 @@ public class DruidAvaticaProtobufHandler extends 
DruidAvaticaHandler
   public boolean handle(Request request, Response response, Callback callback) 
throws Exception
   {
     String requestURI = request.getHttpURI().getPath();
-    if 
(AVATICA_PATH_NO_TRAILING_SLASH.equals(StringUtils.maybeRemoveTrailingSlash(requestURI)))
 {
-      try (Timer.Context ctx = this.requestTimer.start()) {
-        if (!"POST".equals(request.getMethod())) {
-          response.setStatus(405);
-          response.write(
-              true,
-              ByteBuffer.wrap("This server expects only POST 
calls.".getBytes(StandardCharsets.UTF_8)), callback
-          );
-          return true;
-        }
-        final byte[] requestBytes;
-        // Avoid a new buffer creation for every HTTP request
-        final UnsynchronizedBuffer buffer = threadLocalBuffer.get();
-        try (InputStream inputStream = Content.Source.asInputStream(request)) {
-          requestBytes = AvaticaUtils.readFullyToBytes(inputStream, buffer);
-        }
-        finally {
-          buffer.reset();
-        }
+    String remoteAddr = Request.getRemoteAddr(request);
+    DruidMeta.setThreadLocalRemoteAddress(remoteAddr);
 
-        response.getHeaders().put("Content-Type", 
"application/octet-stream;charset=utf-8");
+    try {
+      if 
(AVATICA_PATH_NO_TRAILING_SLASH.equals(StringUtils.maybeRemoveTrailingSlash(requestURI)))
 {
+        try (Timer.Context ctx = this.requestTimer.start()) {
+          if (!"POST".equals(request.getMethod())) {
+            response.setStatus(405);
+            response.write(
+                true,
+                ByteBuffer.wrap("This server expects only POST 
calls.".getBytes(StandardCharsets.UTF_8)), callback
+            );
+            return true;
+          }
+          final byte[] requestBytes;
+          // Avoid a new buffer creation for every HTTP request
+          final UnsynchronizedBuffer buffer = threadLocalBuffer.get();
+          try (InputStream inputStream = 
Content.Source.asInputStream(request)) {
+            requestBytes = AvaticaUtils.readFullyToBytes(inputStream, buffer);
+          }
+          finally {
+            buffer.reset();
+          }
 
-        org.apache.calcite.avatica.remote.Handler.HandlerResponse<byte[]> 
handlerResponse;
-        try {
-          handlerResponse = protobufHandler.apply(requestBytes);
-        }
-        catch (Exception e) {
-          LOG.debug(e, "Error invoking request");
-          handlerResponse = protobufHandler.convertToErrorResponse(e);
-        }
+          response.getHeaders().put("Content-Type", 
"application/octet-stream;charset=utf-8");
+
+          org.apache.calcite.avatica.remote.Handler.HandlerResponse<byte[]> 
handlerResponse;
+          try {
+            handlerResponse = protobufHandler.apply(requestBytes);
+          }
+          catch (Exception e) {
+            LOG.debug(e, "Error invoking request");
+            handlerResponse = protobufHandler.convertToErrorResponse(e);
+          }
 
-        response.setStatus(handlerResponse.getStatusCode());
-        response.write(true, ByteBuffer.wrap(handlerResponse.getResponse()), 
callback);
-        return true;
+          response.setStatus(handlerResponse.getStatusCode());
+          response.write(true, ByteBuffer.wrap(handlerResponse.getResponse()), 
callback);
+          return true;
+        }
       }
+      return false;
+    }
+    finally {
+      DruidMeta.clearThreadLocalRemoteAddress();
     }
-    return false;
   }
 
   @Override
diff --git 
a/sql/src/main/java/org/apache/druid/sql/avatica/DruidConnection.java 
b/sql/src/main/java/org/apache/druid/sql/avatica/DruidConnection.java
index 765839b3fc6..8e8e5cf990a 100644
--- a/sql/src/main/java/org/apache/druid/sql/avatica/DruidConnection.java
+++ b/sql/src/main/java/org/apache/druid/sql/avatica/DruidConnection.java
@@ -139,7 +139,8 @@ public class DruidConnection
       final SqlQueryPlus sqlQueryPlus,
       final Map<String, Object> systemDefaultContext,
       final long maxRowCount,
-      final ResultFetcherFactory fetcherFactory
+      final ResultFetcherFactory fetcherFactory,
+      final String remoteAddress
   )
   {
     final int statementId = statementCounter.incrementAndGet();
@@ -157,7 +158,8 @@ public class DruidConnection
 
       @SuppressWarnings("GuardedBy")
       final PreparedStatement statement = 
sqlStatementFactory.preparedStatement(
-          sqlQueryPlus.withContext(systemDefaultContext, sessionContext)
+          sqlQueryPlus.withContext(systemDefaultContext, sessionContext),
+          remoteAddress
       );
       final DruidJdbcPreparedStatement jdbcStmt = new 
DruidJdbcPreparedStatement(
           connectionId,
diff --git 
a/sql/src/main/java/org/apache/druid/sql/avatica/DruidJdbcPreparedStatement.java
 
b/sql/src/main/java/org/apache/druid/sql/avatica/DruidJdbcPreparedStatement.java
index dcd599c5428..8379512b1e6 100644
--- 
a/sql/src/main/java/org/apache/druid/sql/avatica/DruidJdbcPreparedStatement.java
+++ 
b/sql/src/main/java/org/apache/druid/sql/avatica/DruidJdbcPreparedStatement.java
@@ -29,6 +29,7 @@ import org.apache.druid.sql.PreparedStatement;
 import org.apache.druid.sql.avatica.DruidJdbcResultSet.ResultFetcherFactory;
 import org.apache.druid.sql.calcite.planner.PrepareResult;
 
+import javax.annotation.Nullable;
 import java.util.List;
 
 /**
@@ -94,12 +95,12 @@ public class DruidJdbcPreparedStatement extends 
AbstractDruidJdbcStatement
     return signature;
   }
 
-  public synchronized void execute(List<TypedValue> parameters)
+  public synchronized void execute(List<TypedValue> parameters, @Nullable 
String remoteAddress)
   {
     ensure(State.PREPARED);
     closeResultSet();
     try {
-      DirectStatement directStmt = sqlStatement.execute(parameters);
+      DirectStatement directStmt = sqlStatement.execute(parameters, 
remoteAddress);
       resultSet = new DruidJdbcResultSet(this, directStmt, maxRowCount, 
fetcherFactory);
       resultSet.execute();
     }
diff --git 
a/sql/src/main/java/org/apache/druid/sql/avatica/DruidJdbcStatement.java 
b/sql/src/main/java/org/apache/druid/sql/avatica/DruidJdbcStatement.java
index 3b3a3247f51..6e4570e4a58 100644
--- a/sql/src/main/java/org/apache/druid/sql/avatica/DruidJdbcStatement.java
+++ b/sql/src/main/java/org/apache/druid/sql/avatica/DruidJdbcStatement.java
@@ -26,6 +26,7 @@ import org.apache.druid.sql.SqlQueryPlus;
 import org.apache.druid.sql.SqlStatementFactory;
 import org.apache.druid.sql.avatica.DruidJdbcResultSet.ResultFetcherFactory;
 
+import javax.annotation.Nullable;
 import java.util.Map;
 
 /**
@@ -57,11 +58,17 @@ public class DruidJdbcStatement extends 
AbstractDruidJdbcStatement
     this.lifecycleFactory = Preconditions.checkNotNull(lifecycleFactory, 
"lifecycleFactory");
   }
 
-  public synchronized void execute(SqlQueryPlus queryPlus, long maxRowCount)
+  /**
+   * Executes the given query. The {@code remoteAddress} is read from the 
calling
+   * (execute) request rather than captured when the statement was created, so 
that
+   * a statement reused after a reconnect or from a different remote is 
attributed
+   * to the actual caller of this execution.
+   */
+  public synchronized void execute(SqlQueryPlus queryPlus, long maxRowCount, 
@Nullable String remoteAddress)
   {
     closeResultSet();
     this.sqlQuery = queryPlus.withContext(defaultContext, 
queryContext).freshCopy();
-    DirectStatement stmt = lifecycleFactory.directStatement(this.sqlQuery);
+    DirectStatement stmt = lifecycleFactory.directStatement(this.sqlQuery, 
remoteAddress);
     resultSet = new DruidJdbcResultSet(this, stmt, Long.MAX_VALUE, 
fetcherFactory);
     try {
       resultSet.execute();
diff --git a/sql/src/main/java/org/apache/druid/sql/avatica/DruidMeta.java 
b/sql/src/main/java/org/apache/druid/sql/avatica/DruidMeta.java
index ff27f020832..74524b71eef 100644
--- a/sql/src/main/java/org/apache/druid/sql/avatica/DruidMeta.java
+++ b/sql/src/main/java/org/apache/druid/sql/avatica/DruidMeta.java
@@ -119,6 +119,18 @@ public class DruidMeta extends MetaImpl
 
   private static final Logger LOG = new Logger(DruidMeta.class);
 
+  private static final ThreadLocal<String> THREAD_LOCAL_REMOTE_ADDRESS = new 
ThreadLocal<>();
+
+  public static void setThreadLocalRemoteAddress(String remoteAddress)
+  {
+    THREAD_LOCAL_REMOTE_ADDRESS.set(remoteAddress);
+  }
+
+  public static void clearThreadLocalRemoteAddress()
+  {
+    THREAD_LOCAL_REMOTE_ADDRESS.remove();
+  }
+
   /**
    * Items passed in via the connection context which are not query
    * context values. Instead, these are used at connection time to validate
@@ -267,7 +279,11 @@ public class DruidMeta extends MetaImpl
   {
     try {
       final DruidJdbcStatement druidStatement = getDruidConnection(ch.id)
-          .createStatement(sqlStatementFactory, 
queryConfigProvider.getContext(), fetcherFactory);
+          .createStatement(
+              sqlStatementFactory,
+              queryConfigProvider.getContext(),
+              fetcherFactory
+          );
       return new StatementHandle(ch.id, druidStatement.getStatementId(), null);
     }
     catch (Throwable t) {
@@ -300,7 +316,8 @@ public class DruidMeta extends MetaImpl
           sqlReq,
           queryConfigProvider.getContext(),
           maxRowCount,
-          fetcherFactory
+          fetcherFactory,
+          THREAD_LOCAL_REMOTE_ADDRESS.get()
       );
       stmt.prepare();
       LOG.debug("Successfully prepared statement [%s] for execution", 
stmt.getStatementId());
@@ -371,7 +388,7 @@ public class DruidMeta extends MetaImpl
         final SqlQueryPlus sqlRequest = SqlQueryPlus.builder(sql)
                                                     .auth(authenticationResult)
                                                     .buildJdbc();
-        druidStatement.execute(sqlRequest, maxRowCount);
+        druidStatement.execute(sqlRequest, maxRowCount, 
THREAD_LOCAL_REMOTE_ADDRESS.get());
         final ExecuteResult result = doFetch(druidStatement, 
maxRowsInFirstFrame);
         LOG.debug("Successfully prepared statement [%s] and started 
execution", druidStatement.getStatementId());
         return result;
@@ -494,7 +511,7 @@ public class DruidMeta extends MetaImpl
     try {
       final DruidJdbcPreparedStatement druidStatement =
           getDruidStatement(statement, DruidJdbcPreparedStatement.class);
-      druidStatement.execute(parameterValues);
+      druidStatement.execute(parameterValues, 
THREAD_LOCAL_REMOTE_ADDRESS.get());
       ExecuteResult result = doFetch(druidStatement, maxRowsInFirstFrame);
       LOG.debug(
           "Successfully started execution of statement [%s]",
diff --git a/sql/src/test/java/org/apache/druid/sql/SqlStatementTest.java 
b/sql/src/test/java/org/apache/druid/sql/SqlStatementTest.java
index c4ecd97f507..8856a19331f 100644
--- a/sql/src/test/java/org/apache/druid/sql/SqlStatementTest.java
+++ b/sql/src/test/java/org/apache/druid/sql/SqlStatementTest.java
@@ -436,7 +436,7 @@ public class SqlStatementTest
     // JDBC supports a prepare once, execute many model
     for (int i = 0; i < 3; i++) {
       List<Object[]> results = stmt
-          .execute(Collections.emptyList())
+          .execute(Collections.emptyList(), null)
           .execute()
           .getResults()
           .toList();
@@ -495,7 +495,7 @@ public class SqlStatementTest
         CalciteTests.REGULAR_USER_AUTH_RESULT
     );
     PreparedStatement stmt = sqlStatementFactory.preparedStatement(sqlReq);
-    DruidException e = Assert.assertThrows(DruidException.class, () -> 
stmt.execute(Collections.emptyList()).execute());
+    DruidException e = Assert.assertThrows(DruidException.class, () -> 
stmt.execute(Collections.emptyList(), null).execute());
 
     Assert.assertEquals(DruidException.Category.FORBIDDEN, e.getCategory());
     Assert.assertEquals(DruidException.Persona.OPERATOR, e.getTargetPersona());
@@ -512,7 +512,7 @@ public class SqlStatementTest
         CalciteTests.REGULAR_USER_AUTH_RESULT
     );
     PreparedStatement stmt = sqlStatementFactory.preparedStatement(sqlReq);
-    List<Object[]> results = 
stmt.execute(Collections.emptyList()).execute().getResults().toList();
+    List<Object[]> results = stmt.execute(Collections.emptyList(), 
null).execute().getResults().toList();
 
     ImmutableList<Object[]> expectedResults = ImmutableList.of(new 
Object[]{1L});
     assertResultsEquals("SELECT COUNT(*) AS cnt FROM 
druid.restrictedDatasource_m1_is_6", expectedResults, results);
diff --git 
a/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaHandlerTest.java 
b/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaHandlerTest.java
index 2c963468ef8..20c46f3a108 100644
--- 
a/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaHandlerTest.java
+++ 
b/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaHandlerTest.java
@@ -137,6 +137,7 @@ import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
+import java.util.regex.Pattern;
 
 /**
  * Tests the Avatica-based JDBC implementation using JSON serialization. See
@@ -146,6 +147,9 @@ import java.util.concurrent.TimeUnit;
  */
 public class DruidAvaticaHandlerTest extends CalciteTestBase
 {
+  private static final Pattern IPV4_PATTERN = 
Pattern.compile("^\\d+\\.\\d+\\.\\d+\\.\\d+$");
+  private static final Pattern IPV6_PATTERN = 
Pattern.compile("^[0-9a-fA-F:]+$");
+
   private static final int CONNECTION_LIMIT = 4;
   private static final int STATEMENT_LIMIT = 4;
 
@@ -1950,4 +1954,89 @@ public class DruidAvaticaHandlerTest extends 
CalciteTestBase
     }
     return m;
   }
+
+  /**
+   * Test that remote address is properly captured and logged for JDBC Avatica 
connections.
+   * This verifies the fix for issue #19230 which ensures that the client's 
remote address
+   * is tracked through the entire SQL execution lifecycle.
+   */
+  @Test
+  public void testRemoteAddressInLogs() throws SQLException
+  {
+    testRequestLogger.clear();
+
+    try (Statement stmt = client.createStatement()) {
+      stmt.executeQuery("SELECT COUNT(*) AS cnt FROM druid.foo");
+    }
+
+    Assert.assertEquals(1, testRequestLogger.getSqlQueryLogs().size());
+    RequestLogLine logLine = testRequestLogger.getSqlQueryLogs().get(0);
+
+    String remoteAddress = logLine.getRemoteAddr();
+    Assert.assertNotNull("Remote address should not be null", remoteAddress);
+
+    Assert.assertTrue(
+        "Remote address should be a valid IP address, got: " + remoteAddress,
+        IPV4_PATTERN.matcher(remoteAddress).matches() ||
+        IPV6_PATTERN.matcher(remoteAddress).matches()
+    );
+  }
+
+  /**
+   * Test that remote address is captured even when a query fails.
+   */
+  @Test
+  public void testRemoteAddressInFailedQuery() throws SQLException
+  {
+    testRequestLogger.clear();
+
+    try (Statement stmt = client.createStatement()) {
+      stmt.executeQuery("SELECT nonexistent FROM druid.foo");
+      Assert.fail("Query should have failed");
+    }
+    catch (SQLException e) {
+      // Expected exception
+    }
+
+    Assert.assertEquals(1, testRequestLogger.getSqlQueryLogs().size());
+    RequestLogLine logLine = testRequestLogger.getSqlQueryLogs().get(0);
+
+    String remoteAddress = logLine.getRemoteAddr();
+    Assert.assertNotNull("Remote address should not be null even in failed 
query", remoteAddress);
+    Assert.assertFalse("Remote address should not be empty even in failed 
query", remoteAddress.length() == 0);
+  }
+
+  /**
+   * Test that remote address is captured for prepared statements.
+   * Both the prepare-phase and execute-phase log entries must carry the 
address —
+   * DruidJdbcPreparedStatement.close() emits the prepare-phase reporter, so a
+   * missing address there would leak an empty remoteAddress dimension to 
metrics.
+   */
+  @Test
+  public void testRemoteAddressInPreparedStatement() throws SQLException
+  {
+    testRequestLogger.clear();
+
+    try (PreparedStatement stmt = client.prepareStatement("SELECT COUNT(*) AS 
cnt FROM druid.foo WHERE dim1 = ?")) {
+      stmt.setString(1, "abc");
+      stmt.executeQuery();
+    }
+
+    Assert.assertFalse(
+        "Should have at least one log entry",
+        testRequestLogger.getSqlQueryLogs().isEmpty()
+    );
+
+    for (RequestLogLine logLine : testRequestLogger.getSqlQueryLogs()) {
+      String remoteAddress = logLine.getRemoteAddr();
+      Assert.assertNotNull(
+          "Every prepared-statement log entry must carry a remote address",
+          remoteAddress
+      );
+      Assert.assertFalse(
+          "Every prepared-statement log entry must carry a non-empty remote 
address",
+          remoteAddress.isEmpty()
+      );
+    }
+  }
 }
diff --git 
a/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandlerTest.java
 
b/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandlerTest.java
index 94bb835e0e9..f4a1586c899 100644
--- 
a/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandlerTest.java
+++ 
b/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaJsonHandlerTest.java
@@ -22,12 +22,14 @@ package org.apache.druid.sql.avatica;
 import org.apache.druid.server.DruidNode;
 import org.easymock.EasyMock;
 import org.eclipse.jetty.http.HttpURI;
+import org.eclipse.jetty.server.ConnectionMetaData;
 import org.eclipse.jetty.server.Request;
 import org.eclipse.jetty.server.Response;
 import org.eclipse.jetty.util.Callback;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.net.InetSocketAddress;
 import java.nio.ByteBuffer;
 
 public class DruidAvaticaJsonHandlerTest extends DruidAvaticaHandlerTest
@@ -62,7 +64,11 @@ public class DruidAvaticaJsonHandlerTest extends 
DruidAvaticaHandlerTest
     Response response = EasyMock.mock(Response.class);
     Callback callback = EasyMock.mock(Callback.class);
     HttpURI httpURI = EasyMock.mock(HttpURI.class);
+    ConnectionMetaData connectionMetaData = 
EasyMock.mock(ConnectionMetaData.class);
 
+    
EasyMock.expect(request.getConnectionMetaData()).andReturn(connectionMetaData);
+    EasyMock.expect(connectionMetaData.getRemoteSocketAddress())
+            .andReturn(new InetSocketAddress("127.0.0.1", 12345));
     EasyMock.expect(request.getHttpURI()).andReturn(httpURI);
     
EasyMock.expect(httpURI.getPath()).andReturn(DruidAvaticaProtobufHandler.AVATICA_PATH_NO_TRAILING_SLASH);
     EasyMock.expect(request.getMethod()).andReturn("GET");
@@ -77,11 +83,11 @@ public class DruidAvaticaJsonHandlerTest extends 
DruidAvaticaHandlerTest
     );
     EasyMock.expectLastCall();
 
-    EasyMock.replay(request, response, callback, httpURI);
+    EasyMock.replay(request, response, callback, httpURI, connectionMetaData);
 
     boolean handled = handler.handle(request, response, callback);
 
     Assertions.assertTrue(handled, "Handler should have handled the request");
-    EasyMock.verify(request, response, callback, httpURI);
+    EasyMock.verify(request, response, callback, httpURI, connectionMetaData);
   }
 }
diff --git 
a/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandlerTest.java
 
b/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandlerTest.java
index decf79ed2e0..2eb9ff3bc60 100644
--- 
a/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandlerTest.java
+++ 
b/sql/src/test/java/org/apache/druid/sql/avatica/DruidAvaticaProtobufHandlerTest.java
@@ -23,12 +23,14 @@ import org.apache.druid.java.util.common.StringUtils;
 import org.apache.druid.server.DruidNode;
 import org.easymock.EasyMock;
 import org.eclipse.jetty.http.HttpURI;
+import org.eclipse.jetty.server.ConnectionMetaData;
 import org.eclipse.jetty.server.Request;
 import org.eclipse.jetty.server.Response;
 import org.eclipse.jetty.util.Callback;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.net.InetSocketAddress;
 import java.nio.ByteBuffer;
 
 public class DruidAvaticaProtobufHandlerTest extends DruidAvaticaHandlerTest
@@ -66,26 +68,30 @@ public class DruidAvaticaProtobufHandlerTest extends 
DruidAvaticaHandlerTest
     Response response = EasyMock.mock(Response.class);
     Callback callback = EasyMock.mock(Callback.class);
     HttpURI httpURI = EasyMock.mock(HttpURI.class);
+    ConnectionMetaData connectionMetaData = 
EasyMock.mock(ConnectionMetaData.class);
 
+    
EasyMock.expect(request.getConnectionMetaData()).andReturn(connectionMetaData);
+    EasyMock.expect(connectionMetaData.getRemoteSocketAddress())
+            .andReturn(new InetSocketAddress("127.0.0.1", 12345));
     EasyMock.expect(request.getHttpURI()).andReturn(httpURI);
     
EasyMock.expect(httpURI.getPath()).andReturn(DruidAvaticaProtobufHandler.AVATICA_PATH_NO_TRAILING_SLASH);
     EasyMock.expect(request.getMethod()).andReturn("GET");
-    
+
     response.setStatus(405);
     EasyMock.expectLastCall();
-    
+
     response.write(
-        EasyMock.eq(true), 
-        EasyMock.anyObject(ByteBuffer.class), 
+        EasyMock.eq(true),
+        EasyMock.anyObject(ByteBuffer.class),
         EasyMock.eq(callback)
     );
     EasyMock.expectLastCall();
 
-    EasyMock.replay(request, response, callback, httpURI);
+    EasyMock.replay(request, response, callback, httpURI, connectionMetaData);
 
     boolean handled = handler.handle(request, response, callback);
 
     Assertions.assertTrue(handled, "Handler should have handled the request");
-    EasyMock.verify(request, response, callback, httpURI);
+    EasyMock.verify(request, response, callback, httpURI, connectionMetaData);
   }
 }
diff --git 
a/sql/src/test/java/org/apache/druid/sql/avatica/DruidStatementTest.java 
b/sql/src/test/java/org/apache/druid/sql/avatica/DruidStatementTest.java
index baa7732d1b0..8746a78f33e 100644
--- a/sql/src/test/java/org/apache/druid/sql/avatica/DruidStatementTest.java
+++ b/sql/src/test/java/org/apache/druid/sql/avatica/DruidStatementTest.java
@@ -160,7 +160,7 @@ public class DruidStatementTest extends CalciteTestBase
                     .buildJdbc();
     try (final DruidJdbcStatement statement = jdbcStatement()) {
       // First frame, ask for all rows.
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           subQueryWithOrderByResults(),
@@ -180,7 +180,7 @@ public class DruidStatementTest extends CalciteTestBase
                     .buildJdbc();
     try (final DruidJdbcStatement statement = jdbcStatement()) {
       // First frame, ask for all rows.
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           subQueryWithOrderByResults(),
@@ -227,7 +227,7 @@ public class DruidStatementTest extends CalciteTestBase
                     .buildJdbc();
     try (final DruidJdbcStatement statement = jdbcStatement()) {
       // First frame, ask for all rows.
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       statement.closeResultSet();
       statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
@@ -247,7 +247,7 @@ public class DruidStatementTest extends CalciteTestBase
                     .auth(AllowAllAuthenticator.ALLOW_ALL_RESULT)
                     .buildJdbc();
     try (final DruidJdbcStatement statement = jdbcStatement()) {
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           subQueryWithOrderByResults(),
@@ -256,7 +256,7 @@ public class DruidStatementTest extends CalciteTestBase
 
       // Do it again. JDBC says we can reuse statements sequentially.
       Assert.assertTrue(statement.isDone());
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       frame = statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           subQueryWithOrderByResults(),
@@ -292,7 +292,7 @@ public class DruidStatementTest extends CalciteTestBase
                     .buildJdbc();
     try (final DruidJdbcStatement statement = jdbcStatement()) {
       // First frame, ask for all rows.
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           Meta.Frame.create(
@@ -334,7 +334,7 @@ public class DruidStatementTest extends CalciteTestBase
     try (final DruidJdbcStatement statement = jdbcStatement()) {
 
       // First frame, ask for 2 rows.
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       Assert.assertEquals(0, statement.getCurrentOffset());
       Assert.assertFalse(statement.isDone());
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 2);
@@ -369,7 +369,7 @@ public class DruidStatementTest extends CalciteTestBase
                     .buildJdbc();
     try (final DruidJdbcStatement statement = jdbcStatement()) {
       // First frame, ask for 2 rows.
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 2);
       Assert.assertEquals(
           firstFrameResults(),
@@ -378,7 +378,7 @@ public class DruidStatementTest extends CalciteTestBase
       Assert.assertFalse(statement.isDone());
 
       // Do it again. Closes the prior result set.
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       frame = statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 2);
       Assert.assertEquals(
           firstFrameResults(),
@@ -410,7 +410,7 @@ public class DruidStatementTest extends CalciteTestBase
                     .buildJdbc();
     try (final DruidJdbcStatement statement = jdbcStatement()) {
       // First frame, ask for 2 rows.
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 2);
       Assert.assertEquals(
           firstFrameResults(),
@@ -464,7 +464,7 @@ public class DruidStatementTest extends CalciteTestBase
                     .buildJdbc();
     try (final DruidJdbcStatement statement = jdbcStatement()) {
       // Check signature.
-      statement.execute(queryPlus, -1);
+      statement.execute(queryPlus, -1, null);
       verifySignature(statement.getSignature());
     }
   }
@@ -533,7 +533,7 @@ public class DruidStatementTest extends CalciteTestBase
     try (final DruidJdbcPreparedStatement statement = 
jdbcPreparedStatement(queryPlus)) {
       statement.prepare();
       // First frame, ask for all rows.
-      statement.execute(Collections.emptyList());
+      statement.execute(Collections.emptyList(), null);
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           subQueryWithOrderByResults(),
@@ -554,7 +554,7 @@ public class DruidStatementTest extends CalciteTestBase
                     .buildJdbc();
     try (final DruidJdbcPreparedStatement statement = 
jdbcPreparedStatement(queryPlus)) {
       statement.prepare();
-      statement.execute(Collections.emptyList());
+      statement.execute(Collections.emptyList(), null);
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           subQueryWithOrderByResults(),
@@ -563,7 +563,7 @@ public class DruidStatementTest extends CalciteTestBase
 
       // Do it again. JDBC says we can reuse prepared statements sequentially.
       Assert.assertTrue(statement.isDone());
-      statement.execute(Collections.emptyList());
+      statement.execute(Collections.emptyList(), null);
       frame = statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           subQueryWithOrderByResults(),
@@ -603,7 +603,7 @@ public class DruidStatementTest extends CalciteTestBase
       statement.prepare();
 
       // Execute many times. First time.
-      statement.execute(matchingParams);
+      statement.execute(matchingParams, null);
       Meta.Frame frame = 
statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           expected,
@@ -611,7 +611,7 @@ public class DruidStatementTest extends CalciteTestBase
       );
 
       // Again, same value.
-      statement.execute(matchingParams);
+      statement.execute(matchingParams, null);
       frame = statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           expected,
@@ -621,7 +621,7 @@ public class DruidStatementTest extends CalciteTestBase
       // Again, no matches.
       statement.execute(
           Collections.singletonList(
-              TypedValue.ofLocal(ColumnMetaData.Rep.STRING, "foo")));
+              TypedValue.ofLocal(ColumnMetaData.Rep.STRING, "foo")), null);
       frame = statement.nextFrame(AbstractDruidJdbcStatement.START_OFFSET, 6);
       Assert.assertEquals(
           Meta.Frame.create(0, true, Collections.emptyList()),
diff --git 
a/sql/src/test/java/org/apache/druid/sql/calcite/util/QueryFrameworkUtils.java 
b/sql/src/test/java/org/apache/druid/sql/calcite/util/QueryFrameworkUtils.java
index 2ba1238cd2b..b48743a6d25 100644
--- 
a/sql/src/test/java/org/apache/druid/sql/calcite/util/QueryFrameworkUtils.java
+++ 
b/sql/src/test/java/org/apache/druid/sql/calcite/util/QueryFrameworkUtils.java
@@ -360,7 +360,7 @@ public class QueryFrameworkUtils
     public DirectStatement directStatement(SqlQueryPlus sqlRequest)
     {
       // override direct statement creation to allow calcite tests to test 
multi-part set statements
-      return new DirectStatement(toolbox, sqlRequest)
+      return new DirectStatement(toolbox, sqlRequest, null)
       {
         @Override
         protected DruidPlanner createPlanner()
@@ -380,7 +380,7 @@ public class QueryFrameworkUtils
     @Override
     public PreparedStatement preparedStatement(SqlQueryPlus sqlRequest)
     {
-      return new PreparedStatement(toolbox, sqlRequest)
+      return new PreparedStatement(toolbox, sqlRequest, null)
       {
         @Override
         protected DruidPlanner getPlanner()
@@ -396,7 +396,7 @@ public class QueryFrameworkUtils
         }
 
         @Override
-        public DirectStatement execute(List<TypedValue> parameters)
+        public DirectStatement execute(List<TypedValue> parameters, String 
remoteAddress)
         {
           return directStatement(queryPlus.withParameters(parameters));
         }


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


Reply via email to