From 0bc22f4a3b589bd1e2cbb26094ca5dadc113609f Mon Sep 17 00:00:00 2001 From: Amanda Villarreal Date: Tue, 30 Jun 2026 10:51:36 -0500 Subject: [PATCH 1/8] Removing CustomThreadedSelectorServer.java --- .../rpc/CustomThreadedSelectorServer.java | 74 ------------------- .../accumulo/server/rpc/TServerUtils.java | 2 +- 2 files changed, 1 insertion(+), 75 deletions(-) delete mode 100644 server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java deleted file mode 100644 index 455640ab543..00000000000 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java +++ /dev/null @@ -1,74 +0,0 @@ -/* - * 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 - * - * https://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.accumulo.server.rpc; - -import java.lang.reflect.Field; -import java.net.Socket; - -import org.apache.thrift.server.TThreadedSelectorServer; -import org.apache.thrift.transport.TNonblockingSocket; -import org.apache.thrift.transport.TNonblockingTransport; -import org.slf4j.LoggerFactory; - -public class CustomThreadedSelectorServer extends TThreadedSelectorServer { - - private final Field fbTansportField; - - public CustomThreadedSelectorServer(Args args) { - super(args); - - try { - fbTansportField = FrameBuffer.class.getDeclaredField("trans_"); - fbTansportField.setAccessible(true); - } catch (SecurityException | NoSuchFieldException e) { - throw new IllegalStateException("Failed to access required field in Thrift code.", e); - } - } - - private TNonblockingTransport getTransport(FrameBuffer frameBuffer) { - try { - return (TNonblockingTransport) fbTansportField.get(frameBuffer); - } catch (IllegalAccessException e) { - throw new IllegalStateException(e); - } - } - - @Override - protected Runnable getRunnable(FrameBuffer frameBuffer) { - return () -> { - - try { - TNonblockingTransport transport = getTransport(frameBuffer); - - if (transport instanceof TNonblockingSocket tsock) { - // This block of code makes the client address available to the server side code that - // executes a RPC. It is made available for informational purposes. - Socket sock = tsock.getSocketChannel().socket(); - TServerUtils.clientAddress - .set(sock.getInetAddress().getHostAddress() + ":" + sock.getPort()); - } - } catch (Exception e) { - LoggerFactory.getLogger(CustomThreadedSelectorServer.class) - .warn("Failed to get client address ", e); - } - frameBuffer.invoke(); - }; - } - -} diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java index a601eb8a66c..fc454541539 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java @@ -190,7 +190,7 @@ private static ServerAddress createThreadedSelectorServer(HostAndPort address, address = HostAndPort.fromParts(address.getHost(), transport.getPort()); } - return new ServerAddress(new CustomThreadedSelectorServer(options), address); + return new ServerAddress(new TThreadedSelectorServer(options), address); } /** From 709390e6c67612e66cd449af82a4cac7a8d03591 Mon Sep 17 00:00:00 2001 From: Amanda Villarreal Date: Mon, 6 Jul 2026 10:11:26 -0500 Subject: [PATCH 2/8] Adding CustomThreadedSelectorServer.java back --- .../rpc/CustomThreadedSelectorServer.java | 74 +++++++++++++++++++ 1 file changed, 74 insertions(+) create mode 100644 server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java new file mode 100644 index 00000000000..455640ab543 --- /dev/null +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java @@ -0,0 +1,74 @@ +/* + * 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 + * + * https://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.accumulo.server.rpc; + +import java.lang.reflect.Field; +import java.net.Socket; + +import org.apache.thrift.server.TThreadedSelectorServer; +import org.apache.thrift.transport.TNonblockingSocket; +import org.apache.thrift.transport.TNonblockingTransport; +import org.slf4j.LoggerFactory; + +public class CustomThreadedSelectorServer extends TThreadedSelectorServer { + + private final Field fbTansportField; + + public CustomThreadedSelectorServer(Args args) { + super(args); + + try { + fbTansportField = FrameBuffer.class.getDeclaredField("trans_"); + fbTansportField.setAccessible(true); + } catch (SecurityException | NoSuchFieldException e) { + throw new IllegalStateException("Failed to access required field in Thrift code.", e); + } + } + + private TNonblockingTransport getTransport(FrameBuffer frameBuffer) { + try { + return (TNonblockingTransport) fbTansportField.get(frameBuffer); + } catch (IllegalAccessException e) { + throw new IllegalStateException(e); + } + } + + @Override + protected Runnable getRunnable(FrameBuffer frameBuffer) { + return () -> { + + try { + TNonblockingTransport transport = getTransport(frameBuffer); + + if (transport instanceof TNonblockingSocket tsock) { + // This block of code makes the client address available to the server side code that + // executes a RPC. It is made available for informational purposes. + Socket sock = tsock.getSocketChannel().socket(); + TServerUtils.clientAddress + .set(sock.getInetAddress().getHostAddress() + ":" + sock.getPort()); + } + } catch (Exception e) { + LoggerFactory.getLogger(CustomThreadedSelectorServer.class) + .warn("Failed to get client address ", e); + } + frameBuffer.invoke(); + }; + } + +} From 043854a75d0a0de32c2b3c5060e006f64198dee2 Mon Sep 17 00:00:00 2001 From: Amanda Villarreal Date: Mon, 6 Jul 2026 13:24:20 -0500 Subject: [PATCH 3/8] Removed CustomThreadedServerSelector.java, replaced implemantation with ThriftServerEventHandler.java --- .../rpc/CustomThreadedSelectorServer.java | 74 ------------------- .../accumulo/server/rpc/TServerUtils.java | 6 +- .../server/rpc/ThriftServerEventHandler.java | 66 +++++++++++++++++ 3 files changed, 71 insertions(+), 75 deletions(-) delete mode 100644 server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java create mode 100644 server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java deleted file mode 100644 index 455640ab543..00000000000 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomThreadedSelectorServer.java +++ /dev/null @@ -1,74 +0,0 @@ -/* - * 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 - * - * https://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.accumulo.server.rpc; - -import java.lang.reflect.Field; -import java.net.Socket; - -import org.apache.thrift.server.TThreadedSelectorServer; -import org.apache.thrift.transport.TNonblockingSocket; -import org.apache.thrift.transport.TNonblockingTransport; -import org.slf4j.LoggerFactory; - -public class CustomThreadedSelectorServer extends TThreadedSelectorServer { - - private final Field fbTansportField; - - public CustomThreadedSelectorServer(Args args) { - super(args); - - try { - fbTansportField = FrameBuffer.class.getDeclaredField("trans_"); - fbTansportField.setAccessible(true); - } catch (SecurityException | NoSuchFieldException e) { - throw new IllegalStateException("Failed to access required field in Thrift code.", e); - } - } - - private TNonblockingTransport getTransport(FrameBuffer frameBuffer) { - try { - return (TNonblockingTransport) fbTansportField.get(frameBuffer); - } catch (IllegalAccessException e) { - throw new IllegalStateException(e); - } - } - - @Override - protected Runnable getRunnable(FrameBuffer frameBuffer) { - return () -> { - - try { - TNonblockingTransport transport = getTransport(frameBuffer); - - if (transport instanceof TNonblockingSocket tsock) { - // This block of code makes the client address available to the server side code that - // executes a RPC. It is made available for informational purposes. - Socket sock = tsock.getSocketChannel().socket(); - TServerUtils.clientAddress - .set(sock.getInetAddress().getHostAddress() + ":" + sock.getPort()); - } - } catch (Exception e) { - LoggerFactory.getLogger(CustomThreadedSelectorServer.class) - .warn("Failed to get client address ", e); - } - frameBuffer.invoke(); - }; - } - -} diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java index fc454541539..8ca7b16e20d 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java @@ -190,7 +190,10 @@ private static ServerAddress createThreadedSelectorServer(HostAndPort address, address = HostAndPort.fromParts(address.getHost(), transport.getPort()); } - return new ServerAddress(new TThreadedSelectorServer(options), address); + final TThreadedSelectorServer server = new TThreadedSelectorServer(options); + server.setServerEventHandler(new ThriftServerEventHandler()); + + return new ServerAddress(server, address); } /** @@ -495,6 +498,7 @@ private static ServerAddress createSaslThreadPoolServer(HostAndPort address, TPr final TThreadPoolServer server = createTThreadPoolServer(transport, processor, ugiTransportFactory, protocolFactory, pool); + server.setServerEventHandler(new ThriftServerEventHandler()); return new ServerAddress(server, address); } diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java new file mode 100644 index 00000000000..01499fd3b6b --- /dev/null +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java @@ -0,0 +1,66 @@ +/* + * 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 + * + * https://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.accumulo.server.rpc; + +import java.net.SocketAddress; + +import org.apache.thrift.protocol.TProtocol; +import org.apache.thrift.server.ServerContext; +import org.apache.thrift.server.TServerEventHandler; +import org.apache.thrift.transport.TTransport; + +public class ThriftServerEventHandler implements TServerEventHandler { + + public static class ThriftServerContext implements ServerContext { + + @Override + public T unwrap(Class iface) { + return null; + } + + @Override + public boolean isWrapperFor(Class iface) { + return false; + } + + @Override + public void setRemoteAddress(SocketAddress remoteAddress) { + TServerUtils.clientAddress.set(remoteAddress.toString()); + } + + } + + @Override + public void preServe() {} + + @Override + public ServerContext createContext(TProtocol input, TProtocol output) { + return new ThriftServerContext(); + } + + @Override + public void deleteContext(ServerContext serverContext, TProtocol input, TProtocol output) { + TServerUtils.clientAddress.set(""); + } + + @Override + public void processContext(ServerContext serverContext, TTransport inputTransport, + TTransport outputTransport) {} + +} From 21d9771c87668ab3c3fe3ff4b282900947cb80e6 Mon Sep 17 00:00:00 2001 From: Amanda Villarreal Date: Mon, 6 Jul 2026 14:19:39 -0500 Subject: [PATCH 4/8] Implementing more ThriftServerEventHandler() --- .../apache/accumulo/server/rpc/TServerUtils.java | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java index 8ca7b16e20d..9efb2da1f59 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java @@ -228,7 +228,10 @@ private static ServerAddress createNonBlockingServer(HostAndPort address, TProce address = HostAndPort.fromParts(address.getHost(), transport.getPort()); } - return new ServerAddress(new CustomNonBlockingServer(options), address); + CustomNonBlockingServer server = new CustomNonBlockingServer(options); + server.setServerEventHandler(new ThriftServerEventHandler()); + + return new ServerAddress(server, address); } /** @@ -298,6 +301,8 @@ private static ServerAddress createBlockingServer(HostAndPort address, TProcesso log.info("Blocking Server bound on {}", address); } + server.setServerEventHandler(new ThriftServerEventHandler()); + return new ServerAddress(server, address); } @@ -320,7 +325,11 @@ private static TThreadPoolServer createTThreadPoolServer(TServerTransport transp if (service != null) { options.executorService(service); } - return new TThreadPoolServer(options); + + final TThreadPoolServer server = new TThreadPoolServer(options); + server.setServerEventHandler(new ThriftServerEventHandler()); + + return server; } /** From 40ec5038f11739d77b239b5380d22800e7fe4c91 Mon Sep 17 00:00:00 2001 From: Amanda Villarreal Date: Mon, 6 Jul 2026 14:40:45 -0500 Subject: [PATCH 5/8] 1 more implementation of ThriftServerEventHandler() --- .../java/org/apache/accumulo/server/rpc/TServerUtils.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java index 9efb2da1f59..6359ee00909 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java @@ -407,9 +407,11 @@ private static ServerAddress createSslThreadPoolServer(HostAndPort address, TPro ThreadPoolExecutor pool = createSelfResizingThreadPool(numThreads, threadTimeOut, conf, timeBetweenThreadChecks); + TThreadPoolServer server = createTThreadPoolServer(transport, processor, + ThriftUtil.transportFactory(), protocolFactory, pool); + server.setServerEventHandler(new ThriftServerEventHandler()); - return new ServerAddress(createTThreadPoolServer(transport, processor, - ThriftUtil.transportFactory(), protocolFactory, pool), address); + return new ServerAddress(server, address); } private static ServerAddress createSaslThreadPoolServer(HostAndPort address, TProcessor processor, From fd591fba099fd0f1730aa9092e785d2c6d613d21 Mon Sep 17 00:00:00 2001 From: Amanda Villarreal Date: Thu, 9 Jul 2026 13:22:05 -0500 Subject: [PATCH 6/8] pr updates, not complete --- .../org/apache/accumulo/server/rpc/TServerUtils.java | 7 +++---- .../accumulo/server/rpc/ThriftServerEventHandler.java | 11 ++++++++--- 2 files changed, 11 insertions(+), 7 deletions(-) diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java index 6359ee00909..cee5e1f7d65 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java @@ -51,6 +51,8 @@ import org.apache.thrift.TProcessor; import org.apache.thrift.TProcessorFactory; import org.apache.thrift.protocol.TProtocolFactory; +import org.apache.thrift.server.THsHaServer; +import org.apache.thrift.server.TNonblockingServer; import org.apache.thrift.server.TThreadPoolServer; import org.apache.thrift.server.TThreadedSelectorServer; import org.apache.thrift.transport.TNonblockingServerSocket; @@ -210,7 +212,7 @@ private static ServerAddress createNonBlockingServer(HostAndPort address, TProce .clientTimeout(0).maxFrameSize(Ints.saturatedCast(maxMessageSize)); final TNonblockingServerSocket transport = new TNonblockingServerSocket(args); - final CustomNonBlockingServer.Args options = new CustomNonBlockingServer.Args(transport); + final THsHaServer server = new THsHaServer(); options.protocolFactory(protocolFactory); options.transportFactory(ThriftUtil.transportFactory(maxMessageSize)); @@ -228,9 +230,6 @@ private static ServerAddress createNonBlockingServer(HostAndPort address, TProce address = HostAndPort.fromParts(address.getHost(), transport.getPort()); } - CustomNonBlockingServer server = new CustomNonBlockingServer(options); - server.setServerEventHandler(new ThriftServerEventHandler()); - return new ServerAddress(server, address); } diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java index 01499fd3b6b..325c2d23e21 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java @@ -27,10 +27,12 @@ public class ThriftServerEventHandler implements TServerEventHandler { + private static final ThriftServerContext context = new ThriftServerContext(); + public static class ThriftServerContext implements ServerContext { @Override - public T unwrap(Class iface) { + public T unwrap(Class iface) throws UnsupportedOperationException { return null; } @@ -44,6 +46,9 @@ public void setRemoteAddress(SocketAddress remoteAddress) { TServerUtils.clientAddress.set(remoteAddress.toString()); } + public void clear() { + TServerUtils.clientAddress.set(""); + } } @Override @@ -51,12 +56,12 @@ public void preServe() {} @Override public ServerContext createContext(TProtocol input, TProtocol output) { - return new ThriftServerContext(); + return context; } @Override public void deleteContext(ServerContext serverContext, TProtocol input, TProtocol output) { - TServerUtils.clientAddress.set(""); + context.clear(); } @Override From ab670d8f82b0cd2b564037e9845315b969d22f3d Mon Sep 17 00:00:00 2001 From: Amanda Villarreal Date: Thu, 9 Jul 2026 14:27:12 -0500 Subject: [PATCH 7/8] Deleting CustomNonBlockingServer.java --- .../server/rpc/CustomNonBlockingServer.java | 179 ------------------ .../accumulo/server/rpc/TServerUtils.java | 6 +- 2 files changed, 4 insertions(+), 181 deletions(-) delete mode 100644 server/base/src/main/java/org/apache/accumulo/server/rpc/CustomNonBlockingServer.java diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomNonBlockingServer.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomNonBlockingServer.java deleted file mode 100644 index 6cbd5c6ca52..00000000000 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/CustomNonBlockingServer.java +++ /dev/null @@ -1,179 +0,0 @@ -/* - * 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 - * - * https://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.accumulo.server.rpc; - -import java.io.IOException; -import java.lang.reflect.Field; -import java.net.Socket; -import java.nio.channels.SelectionKey; - -import org.apache.thrift.server.THsHaServer; -import org.apache.thrift.server.TNonblockingServer; -import org.apache.thrift.transport.TNonblockingServerTransport; -import org.apache.thrift.transport.TNonblockingSocket; -import org.apache.thrift.transport.TNonblockingTransport; -import org.apache.thrift.transport.TTransportException; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * This class implements a custom non-blocking thrift server that stores the client address in - * thread-local storage for the invocation. - */ -public class CustomNonBlockingServer extends THsHaServer { - - private static final Logger log = LoggerFactory.getLogger(CustomNonBlockingServer.class); - private final Field selectAcceptThreadField; - - public CustomNonBlockingServer(Args args) { - super(args); - - try { - selectAcceptThreadField = TNonblockingServer.class.getDeclaredField("selectAcceptThread_"); - selectAcceptThreadField.setAccessible(true); - } catch (ReflectiveOperationException e) { - throw new IllegalStateException("Failed to access required field in Thrift code.", e); - } - } - - @Override - public void stop() { - super.stop(); - try { - getInvoker().shutdownNow(); - } catch (Exception e) { - log.error("Unable to call shutdownNow", e); - } - } - - @Override - protected boolean startThreads() { - // Yet another dirty/gross hack to get access to the client's address. - - // start the selector - try { - // Hack in our SelectAcceptThread impl - SelectAcceptThread selectAcceptThread_ = - new CustomSelectAcceptThread((TNonblockingServerTransport) serverTransport_); - // Set the private field before continuing. - selectAcceptThreadField.set(this, selectAcceptThread_); - - selectAcceptThread_.start(); - return true; - } catch (IOException e) { - LOGGER.error("Failed to start selector thread!", e); - return false; - } catch (IllegalAccessException | IllegalArgumentException e) { - throw new IllegalStateException("Exception setting customer select thread in Thrift"); - } - } - - /** - * Custom wrapper around {@link org.apache.thrift.server.TNonblockingServer.SelectAcceptThread} to - * create our {@link CustomFrameBuffer}. - */ - private class CustomSelectAcceptThread extends SelectAcceptThread { - - public CustomSelectAcceptThread(TNonblockingServerTransport serverTransport) - throws IOException { - super(serverTransport); - } - - @Override - protected FrameBuffer createFrameBuffer(final TNonblockingTransport trans, - final SelectionKey selectionKey, final AbstractSelectThread selectThread) - throws TTransportException { - if (processorFactory_.isAsyncProcessor()) { - throw new IllegalStateException("This implementation does not support AsyncProcessors"); - } - - return new CustomFrameBuffer(trans, selectionKey, selectThread); - } - } - - /** - * Custom wrapper around {@link org.apache.thrift.server.AbstractNonblockingServer.FrameBuffer} to - * extract the client's network location before accepting the request. - */ - private class CustomFrameBuffer extends FrameBuffer { - private final String clientAddress; - - public CustomFrameBuffer(TNonblockingTransport trans, SelectionKey selectionKey, - AbstractSelectThread selectThread) throws TTransportException { - super(trans, selectionKey, selectThread); - // Store the clientAddress in the buffer so it can be referenced for logging during read/write - this.clientAddress = getClientAddress(); - } - - @Override - public void invoke() { - // On invoke() set the clientAddress on the ThreadLocal so that it can be accessed elsewhere - // in the same thread that called invoke() on the buffer - TServerUtils.clientAddress.set(clientAddress); - super.invoke(); - } - - @Override - public boolean read() { - boolean result = super.read(); - if (!result) { - log.trace("CustomFrameBuffer.read returned false when reading data from client: {}", - clientAddress); - } - return result; - } - - @Override - public boolean write() { - boolean result = super.write(); - if (!result) { - log.trace("CustomFrameBuffer.write returned false when writing data to client: {}", - clientAddress); - } - return result; - } - - /* - * Helper method used to capture the client address inside the CustomFrameBuffer constructor so - * that it can be referenced inside the read/write methods for logging purposes. It previously - * was only set on the ThreadLocal in the invoke() method but that does not work because A) the - * method isn't called until after reading is finished so the value will be null inside of - * read() and B) The other problem is that invoke() is called on a different thread than - * read()/write() so even if the order was correct it would not be available. - * - * Since a new FrameBuffer is created for each request we can use it to capture the client - * address earlier in the constructor and not wait for invoke(). A FrameBuffer is used to read - * data and write a response back to the client and as part of creation of the buffer the - * TNonblockingSocket is stored as a final variable and won't change so we can safely capture - * the clientAddress in the constructor and use it for logging during read/write and then use - * the value inside of invoke() to set the ThreadLocal so the client address will still be - * available on the thread that called invoke(). - */ - private String getClientAddress() { - String clientAddress = null; - if (trans_ instanceof TNonblockingSocket tsock) { - Socket sock = tsock.getSocketChannel().socket(); - clientAddress = sock.getInetAddress().getHostAddress() + ":" + sock.getPort(); - log.trace("CustomFrameBuffer captured client address: {}", clientAddress); - } - return clientAddress; - } - } - -} diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java index cee5e1f7d65..86fc0d51aee 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java @@ -52,7 +52,6 @@ import org.apache.thrift.TProcessorFactory; import org.apache.thrift.protocol.TProtocolFactory; import org.apache.thrift.server.THsHaServer; -import org.apache.thrift.server.TNonblockingServer; import org.apache.thrift.server.TThreadPoolServer; import org.apache.thrift.server.TThreadedSelectorServer; import org.apache.thrift.transport.TNonblockingServerSocket; @@ -212,7 +211,7 @@ private static ServerAddress createNonBlockingServer(HostAndPort address, TProce .clientTimeout(0).maxFrameSize(Ints.saturatedCast(maxMessageSize)); final TNonblockingServerSocket transport = new TNonblockingServerSocket(args); - final THsHaServer server = new THsHaServer(); + THsHaServer.Args options = new THsHaServer.Args(transport); options.protocolFactory(protocolFactory); options.transportFactory(ThriftUtil.transportFactory(maxMessageSize)); @@ -230,6 +229,9 @@ private static ServerAddress createNonBlockingServer(HostAndPort address, TProce address = HostAndPort.fromParts(address.getHost(), transport.getPort()); } + final THsHaServer server = new THsHaServer(options); + server.setServerEventHandler(new ThriftServerEventHandler()); + return new ServerAddress(server, address); } From 8fb6937f41dc144e07485a09dd775eb03af7aa41 Mon Sep 17 00:00:00 2001 From: Amanda Villarreal Date: Fri, 10 Jul 2026 11:50:30 -0500 Subject: [PATCH 8/8] Deleting ClientInfoProcessorFactory.java and its 1 implementation --- .../rpc/ClientInfoProcessorFactory.java | 51 ------------------- .../accumulo/server/rpc/TServerUtils.java | 4 +- .../server/rpc/ThriftServerEventHandler.java | 4 +- 3 files changed, 3 insertions(+), 56 deletions(-) delete mode 100644 server/base/src/main/java/org/apache/accumulo/server/rpc/ClientInfoProcessorFactory.java diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/ClientInfoProcessorFactory.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/ClientInfoProcessorFactory.java deleted file mode 100644 index 94d0bd71414..00000000000 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/ClientInfoProcessorFactory.java +++ /dev/null @@ -1,51 +0,0 @@ -/* - * 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 - * - * https://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.accumulo.server.rpc; - -import org.apache.thrift.TProcessor; -import org.apache.thrift.TProcessorFactory; -import org.apache.thrift.transport.TSocket; -import org.apache.thrift.transport.TTransport; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * Sets the address of a client in a ThreadLocal to allow for more informative log messages. - */ -public class ClientInfoProcessorFactory extends TProcessorFactory { - private static final Logger log = LoggerFactory.getLogger(ClientInfoProcessorFactory.class); - - private final ThreadLocal clientAddress; - - public ClientInfoProcessorFactory(ThreadLocal clientAddress, TProcessor processor) { - super(processor); - this.clientAddress = clientAddress; - } - - @Override - public TProcessor getProcessor(TTransport trans) { - if (trans instanceof TSocket tsock) { - clientAddress.set( - tsock.getSocket().getInetAddress().getHostAddress() + ":" + tsock.getSocket().getPort()); - } else { - log.warn("Unable to extract clientAddress from transport of type {}", trans.getClass()); - } - return super.getProcessor(trans); - } -} diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java index 86fc0d51aee..67503e57a1a 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/TServerUtils.java @@ -76,8 +76,7 @@ public class TServerUtils { private static final Logger log = LoggerFactory.getLogger(TServerUtils.class); /** - * Static instance, passed to {@link ClientInfoProcessorFactory}, which will contain the client - * address of any incoming RPC. + * Static instance, which will contain the client address of any incoming RPC. */ public static final ThreadLocal clientAddress = new ThreadLocal<>(); @@ -322,7 +321,6 @@ private static TThreadPoolServer createTThreadPoolServer(TServerTransport transp TThreadPoolServer.Args options = new TThreadPoolServer.Args(transport); options.protocolFactory(protocolFactory); options.transportFactory(transportFactory); - options.processorFactory(new ClientInfoProcessorFactory(clientAddress, processor)); if (service != null) { options.executorService(service); } diff --git a/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java b/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java index 325c2d23e21..a4564653c4c 100644 --- a/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java +++ b/server/base/src/main/java/org/apache/accumulo/server/rpc/ThriftServerEventHandler.java @@ -32,8 +32,8 @@ public class ThriftServerEventHandler implements TServerEventHandler { public static class ThriftServerContext implements ServerContext { @Override - public T unwrap(Class iface) throws UnsupportedOperationException { - return null; + public T unwrap(Class iface) { + throw new UnsupportedOperationException("This method has not been implemented"); } @Override