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

vavrtom pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/qpid-broker-j.git


The following commit(s) were added to refs/heads/main by this push:
     new 733a7914ac QPID-8743: [Broker-J] Add operation to close idle 
connections (#400)
733a7914ac is described below

commit 733a7914ac0e415b39f21ad1adff3229a5bf9254
Author: Daniil Kirilyuk <[email protected]>
AuthorDate: Tue Jun 23 09:03:25 2026 +0200

    QPID-8743: [Broker-J] Add operation to close idle connections (#400)
---
 .../server/virtualhost/AbstractVirtualHost.java    | 216 +++++++++++++++++++++
 .../virtualhost/QueueManagingVirtualHost.java      |   6 +
 .../apache/qpid/server/model/VirtualHostTest.java  | 164 ++++++++++++++++
 .../latest/LatestManagementControllerTest.java     |  71 +++++++
 .../qpid/tests/http/rest/model/OperationTest.java  |  72 +++++++
 5 files changed, 529 insertions(+)

diff --git 
a/broker-core/src/main/java/org/apache/qpid/server/virtualhost/AbstractVirtualHost.java
 
b/broker-core/src/main/java/org/apache/qpid/server/virtualhost/AbstractVirtualHost.java
index a7eba97150..c9a448bd35 100644
--- 
a/broker-core/src/main/java/org/apache/qpid/server/virtualhost/AbstractVirtualHost.java
+++ 
b/broker-core/src/main/java/org/apache/qpid/server/virtualhost/AbstractVirtualHost.java
@@ -44,11 +44,15 @@ import java.util.EnumSet;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.Iterator;
+import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Objects;
 import java.util.Set;
 import java.util.UUID;
+import java.util.concurrent.CancellationException;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ConcurrentMap;
 import java.util.concurrent.ExecutionException;
@@ -56,6 +60,7 @@ import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.ScheduledFuture;
 import java.util.concurrent.ScheduledThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicLong;
 import java.util.regex.Pattern;
@@ -164,6 +169,7 @@ import org.apache.qpid.server.txn.DtxRegistry;
 import org.apache.qpid.server.txn.LocalTransaction;
 import org.apache.qpid.server.txn.ServerTransaction;
 import org.apache.qpid.server.util.HousekeepingExecutor;
+import org.apache.qpid.server.util.ServerScopedRuntimeException;
 import org.apache.qpid.server.util.Strings;
 
 public abstract class AbstractVirtualHost<X extends AbstractVirtualHost<X>> 
extends AbstractConfiguredObject<X>
@@ -173,6 +179,10 @@ public abstract class AbstractVirtualHost<X extends 
AbstractVirtualHost<X>> exte
     private static final Logger LOGGER = 
LoggerFactory.getLogger(AbstractVirtualHost.class);
     private static final Logger DIRECT_MEMORY_USAGE_LOGGER = 
LoggerFactory.getLogger("org.apache.qpid.server.directMemory.virtualhost");
     private static final int HOUSEKEEPING_SHUTDOWN_TIMEOUT = 5;
+    private static final String CONNECTION_CLOSE_STATUS_CLOSE_REQUESTED = 
"CLOSE_REQUESTED";
+    private static final String CONNECTION_CLOSE_STATUS_CLOSED = "CLOSED";
+    private static final String CONNECTION_CLOSE_STATUS_FAILED = "FAILED";
+    private static final String CONNECTION_CLOSE_STATUS_TIMED_OUT = 
"TIMED_OUT";
 
     private final Collection<ConnectionValidator> _connectionValidators = new 
ArrayList<>();
     private final Set<AMQPConnection<?>> _connections = newSetFromMap(new 
ConcurrentHashMap<>());
@@ -756,6 +766,212 @@ public abstract class AbstractVirtualHost<X extends 
AbstractVirtualHost<X>> exte
         return null;
     }
 
+    @Override
+    public Map<String, Object> closeIdleConnections(final long idleTimeMillis,
+                                                    final boolean waitForClose,
+                                                    final long timeoutMillis)
+    {
+        if (idleTimeMillis < 0L)
+        {
+            throw new IllegalArgumentException("idleTimeMillis must be greater 
than or equal to zero");
+        }
+        if (timeoutMillis <= 0L)
+        {
+            throw new IllegalArgumentException("timeoutMillis must be greater 
than zero");
+        }
+
+        final long now = System.currentTimeMillis();
+        final List<IdleConnectionCloseResult> closeResults = new ArrayList<>();
+
+        for (final AMQPConnection<?> connection : new 
ArrayList<>(_connections))
+        {
+            final Date lastMessageTime = connection.getLastMessageTime() != 
null
+                    ? connection.getLastMessageTime()
+                    : connection.getCreatedTime();
+            if (lastMessageTime != null)
+            {
+                final long connectionIdleTimeMillis = now - 
lastMessageTime.getTime();
+                if (connectionIdleTimeMillis > idleTimeMillis)
+                {
+                    closeResults.add(closeIdleConnection(connection, 
lastMessageTime, connectionIdleTimeMillis));
+                }
+            }
+        }
+
+        if (waitForClose)
+        {
+            waitForIdleConnectionCloseResults(closeResults, timeoutMillis);
+        }
+
+        return createCloseIdleConnectionsSummary(closeResults);
+    }
+
+    private IdleConnectionCloseResult closeIdleConnection(final 
AMQPConnection<?> connection,
+                                                          final Date 
lastMessageTime,
+                                                          final long 
idleTimeMillis)
+    {
+        final IdleConnectionCloseResult closeResult = new 
IdleConnectionCloseResult(connection, lastMessageTime,
+                idleTimeMillis);
+        try
+        {
+            closeResult.setCloseFuture(connection.deleteAsync());
+        }
+        catch (final RuntimeException e)
+        {
+            closeResult.setFailed(e);
+        }
+        return closeResult;
+    }
+
+    private void waitForIdleConnectionCloseResults(final 
List<IdleConnectionCloseResult> closeResults,
+                                                   final long timeoutMillis)
+    {
+        final CompletableFuture<?>[] closeFutures = closeResults.stream()
+                .map(IdleConnectionCloseResult::getCloseFuture)
+                .filter(Objects::nonNull)
+                .toArray(CompletableFuture[]::new);
+
+        if (closeFutures.length > 0)
+        {
+            try
+            {
+                CompletableFuture.allOf(closeFutures).get(timeoutMillis, 
TimeUnit.MILLISECONDS);
+            }
+            catch (final TimeoutException | ExecutionException e)
+            {
+                // Individual result status is derived below from each close 
future.
+            }
+            catch (final InterruptedException e)
+            {
+                Thread.currentThread().interrupt();
+                throw new ServerScopedRuntimeException("Interrupted while 
waiting for idle connections to close", e);
+            }
+        }
+
+        closeResults.forEach(IdleConnectionCloseResult::updateStatusAfterWait);
+    }
+
+    private Map<String, Object> createCloseIdleConnectionsSummary(final 
List<IdleConnectionCloseResult> closeResults)
+    {
+        int closeRequestedCount = 0;
+        int closedCount = 0;
+        int failedCount = 0;
+        int timedOutCount = 0;
+
+        final List<Map<String, Object>> connections = new 
ArrayList<>(closeResults.size());
+        for (final IdleConnectionCloseResult closeResult : closeResults)
+        {
+            final Map<String, Object> connection = 
closeResult.getConnectionSummary();
+            connections.add(connection);
+
+            final String status = (String) connection.get("status");
+            if (closeResult.getCloseFuture() != null)
+            {
+                closeRequestedCount++;
+            }
+            if (CONNECTION_CLOSE_STATUS_CLOSED.equals(status))
+            {
+                closedCount++;
+            }
+            else if (CONNECTION_CLOSE_STATUS_FAILED.equals(status))
+            {
+                failedCount++;
+            }
+            else if (CONNECTION_CLOSE_STATUS_TIMED_OUT.equals(status))
+            {
+                timedOutCount++;
+            }
+        }
+
+        final Map<String, Object> result = new LinkedHashMap<>();
+        result.put("matchedCount", closeResults.size());
+        result.put("closeRequestedCount", closeRequestedCount);
+        result.put("closedCount", closedCount);
+        result.put("failedCount", failedCount);
+        result.put("timedOutCount", timedOutCount);
+        result.put("connections", connections);
+        return result;
+    }
+
+    private static class IdleConnectionCloseResult
+    {
+        private final Map<String, Object> _connectionSummary;
+        private CompletableFuture<Void> _closeFuture;
+
+        private IdleConnectionCloseResult(final Connection<?> connection,
+                                          final Date lastMessageTime,
+                                          final long idleTimeMillis)
+        {
+            _connectionSummary = new LinkedHashMap<>();
+            _connectionSummary.put("id", connection.getId() == null ? null : 
connection.getId().toString());
+            _connectionSummary.put("name", connection.getName());
+            _connectionSummary.put("remoteAddress", 
connection.getRemoteAddress());
+            _connectionSummary.put("principal", connection.getPrincipal());
+            _connectionSummary.put("lastMessageTime", lastMessageTime == null 
? null : lastMessageTime.toInstant().toString());
+            _connectionSummary.put("idleTimeMillis", idleTimeMillis);
+            _connectionSummary.put("status", 
CONNECTION_CLOSE_STATUS_CLOSE_REQUESTED);
+        }
+
+        private Map<String, Object> getConnectionSummary()
+        {
+            return _connectionSummary;
+        }
+
+        private CompletableFuture<Void> getCloseFuture()
+        {
+            return _closeFuture;
+        }
+
+        private void setCloseFuture(final CompletableFuture<Void> closeFuture)
+        {
+            if (closeFuture == null)
+            {
+                setFailed(new IllegalStateException("Close request did not 
return a future"));
+            }
+            else
+            {
+                _closeFuture = closeFuture;
+            }
+        }
+
+        private void updateStatusAfterWait()
+        {
+            if (_closeFuture != null)
+            {
+                if (_closeFuture.isDone())
+                {
+                    try
+                    {
+                        _closeFuture.getNow(null);
+                        _connectionSummary.put("status", 
CONNECTION_CLOSE_STATUS_CLOSED);
+                    }
+                    catch (CompletionException | CancellationException e)
+                    {
+                        setFailed(e);
+                    }
+                }
+                else
+                {
+                    _connectionSummary.put("status", 
CONNECTION_CLOSE_STATUS_TIMED_OUT);
+                }
+            }
+        }
+
+        private void setFailed(final Throwable e)
+        {
+            _connectionSummary.put("status", CONNECTION_CLOSE_STATUS_FAILED);
+            _connectionSummary.put("failure", getFailureMessage(e));
+        }
+
+        private static String getFailureMessage(final Throwable e)
+        {
+            final Throwable cause = e instanceof CompletionException && 
e.getCause() != null ? e.getCause() : e;
+            final String causeClassName = cause.getClass().getName();
+            final String message = cause.getMessage() == null ? "null" : 
cause.getMessage();
+            return causeClassName + ": " + message;
+        }
+    }
+
     @Override
     public int publishMessage(@Param(name = "message") final ManageableMessage 
message)
     {
diff --git 
a/broker-core/src/main/java/org/apache/qpid/server/virtualhost/QueueManagingVirtualHost.java
 
b/broker-core/src/main/java/org/apache/qpid/server/virtualhost/QueueManagingVirtualHost.java
index 75f61b9175..c1c17e67f2 100644
--- 
a/broker-core/src/main/java/org/apache/qpid/server/virtualhost/QueueManagingVirtualHost.java
+++ 
b/broker-core/src/main/java/org/apache/qpid/server/virtualhost/QueueManagingVirtualHost.java
@@ -321,6 +321,12 @@ public interface QueueManagingVirtualHost<X extends 
QueueManagingVirtualHost<X>>
     @ManagedOperation(nonModifying = true, changesConfiguredObjectState = 
false)
     Connection<?> getConnection(@Param(name="name", mandatory = true) String 
name);
 
+    @ManagedOperation(description = "Close messaging connections idle for 
longer than the supplied interval.",
+            changesConfiguredObjectState = false)
+    Map<String, Object> closeIdleConnections(@Param(name = "idleTimeMillis", 
mandatory = true) long idleTimeMillis,
+                                             @Param(name = "waitForClose", 
defaultValue = "false") boolean waitForClose,
+                                             @Param(name = "timeoutMillis", 
defaultValue = "10000") long timeoutMillis);
+
     @ManagedOperation(secure = true,
             description = "Publishes a message to a specified address. "
                           + "Returns the number of queues onto which it has 
been placed, "
diff --git 
a/broker-core/src/test/java/org/apache/qpid/server/model/VirtualHostTest.java 
b/broker-core/src/test/java/org/apache/qpid/server/model/VirtualHostTest.java
index 6e4233138b..2cff1122e9 100644
--- 
a/broker-core/src/test/java/org/apache/qpid/server/model/VirtualHostTest.java
+++ 
b/broker-core/src/test/java/org/apache/qpid/server/model/VirtualHostTest.java
@@ -464,6 +464,139 @@ public class VirtualHostTest extends UnitTestBase
         assertEquals(Set.of(connection), vhost.getConnections(), "unexpected 
connection object");
     }
 
+    @Test
+    public void testCloseIdleConnectionsClosesOnlyConnectionsBeyondThreshold()
+    {
+        final QueueManagingVirtualHost<?> vhost = 
createVirtualHost(getTestName());
+        final AMQPConnection<?> idleConnection = 
mockAmqpConnection(mockAuthenticatedPrincipal("user1"),
+                "idleConnection", new Date(System.currentTimeMillis() - 
60000L), CompletableFuture.completedFuture(null));
+        final AMQPConnection<?> activeConnection = 
mockAmqpConnection(mockAuthenticatedPrincipal("user2"),
+                "activeConnection", new Date(), 
CompletableFuture.completedFuture(null));
+
+        vhost.registerConnection(idleConnection);
+        vhost.registerConnection(activeConnection);
+
+        final Map<String, Object> result = vhost.closeIdleConnections(30000L, 
false, 2000L);
+
+        assertEquals(1, result.get("matchedCount"), "Unexpected matched 
count");
+        verify(idleConnection).deleteAsync();
+        verify(activeConnection, never()).deleteAsync();
+    }
+
+    @Test
+    public void testCloseIdleConnectionsReturnsTableSummaryForCloseRequested()
+    {
+        final QueueManagingVirtualHost<?> vhost = 
createVirtualHost(getTestName());
+        final Date lastMessageTime = new Date(System.currentTimeMillis() - 
60000L);
+        final CompletableFuture<Void> closeFuture = new CompletableFuture<>();
+        final AMQPConnection<?> connection = 
mockAmqpConnection(mockAuthenticatedPrincipal("user1"),
+                "connection1", lastMessageTime, closeFuture);
+        vhost.registerConnection(connection);
+
+        final Map<String, Object> result = vhost.closeIdleConnections(0L, 
false, 2000L);
+
+        assertEquals(1, result.get("matchedCount"), "Unexpected matched 
count");
+        assertEquals(1, result.get("closeRequestedCount"), "Unexpected close 
requested count");
+        assertEquals(0, result.get("closedCount"), "Unexpected closed count");
+        assertEquals(0, result.get("failedCount"), "Unexpected failed count");
+        assertEquals(0, result.get("timedOutCount"), "Unexpected timed out 
count");
+
+        final Map<String, Object> connectionResult = 
getSingleConnectionResult(result);
+        assertEquals(connection.getId().toString(), 
connectionResult.get("id"), "Unexpected connection id");
+        assertEquals("connection1", connectionResult.get("name"), "Unexpected 
connection name");
+        assertEquals("/127.0.0.1:5672", connectionResult.get("remoteAddress"), 
"Unexpected remote address");
+        assertEquals("user1", connectionResult.get("principal"), "Unexpected 
principal");
+        assertEquals(lastMessageTime.toInstant().toString(), 
connectionResult.get("lastMessageTime"),
+                "Unexpected last message time");
+        assertTrue(((Number) 
connectionResult.get("idleTimeMillis")).longValue() >= 60000L,
+                "Unexpected idle time");
+        assertEquals("CLOSE_REQUESTED", connectionResult.get("status"), 
"Unexpected status");
+    }
+
+    @Test
+    public void testCloseIdleConnectionsWaitsForCompletedClose()
+    {
+        final QueueManagingVirtualHost<?> vhost = 
createVirtualHost(getTestName());
+        final AMQPConnection<?> connection = 
mockAmqpConnection(mockAuthenticatedPrincipal("user1"),
+                "connection1", new Date(System.currentTimeMillis() - 60000L), 
CompletableFuture.completedFuture(null));
+        vhost.registerConnection(connection);
+
+        final Map<String, Object> result = vhost.closeIdleConnections(0L, 
true, 2000L);
+
+        assertEquals(1, result.get("matchedCount"), "Unexpected matched 
count");
+        assertEquals(1, result.get("closeRequestedCount"), "Unexpected close 
requested count");
+        assertEquals(1, result.get("closedCount"), "Unexpected closed count");
+        assertEquals(0, result.get("failedCount"), "Unexpected failed count");
+        assertEquals(0, result.get("timedOutCount"), "Unexpected timed out 
count");
+        assertEquals("CLOSED", 
getSingleConnectionResult(result).get("status"), "Unexpected status");
+    }
+
+    @Test
+    public void testCloseIdleConnectionsWaitsForTotalTimeoutBudget()
+    {
+        final QueueManagingVirtualHost<?> vhost = 
createVirtualHost(getTestName());
+        final CompletableFuture<Void> closeFuture = new CompletableFuture<>();
+        final AMQPConnection<?> connection = 
mockAmqpConnection(mockAuthenticatedPrincipal("user1"),
+                "connection1", new Date(System.currentTimeMillis() - 60000L), 
closeFuture);
+        vhost.registerConnection(connection);
+
+        final Map<String, Object> result = vhost.closeIdleConnections(0L, 
true, 1L);
+
+        assertEquals(1, result.get("matchedCount"), "Unexpected matched 
count");
+        assertEquals(1, result.get("closeRequestedCount"), "Unexpected close 
requested count");
+        assertEquals(0, result.get("closedCount"), "Unexpected closed count");
+        assertEquals(0, result.get("failedCount"), "Unexpected failed count");
+        assertEquals(1, result.get("timedOutCount"), "Unexpected timed out 
count");
+        assertEquals("TIMED_OUT", 
getSingleConnectionResult(result).get("status"), "Unexpected status");
+        assertFalse(closeFuture.isCancelled(), "Close future should not be 
cancelled");
+    }
+
+    @Test
+    public void testCloseIdleConnectionsReportsFailures()
+    {
+        final QueueManagingVirtualHost<?> vhost = 
createVirtualHost(getTestName());
+        final AMQPConnection<?> thrownFailure = 
mockAmqpConnection(mockAuthenticatedPrincipal("user1"),
+                "thrownFailure", new Date(System.currentTimeMillis() - 
60000L), CompletableFuture.completedFuture(null));
+        when(thrownFailure.deleteAsync()).thenThrow(new 
IllegalStateException("request failed"));
+
+        final CompletableFuture<Void> failedFuture = new CompletableFuture<>();
+        failedFuture.completeExceptionally(new IllegalStateException("future 
failed"));
+        final AMQPConnection<?> futureFailure = 
mockAmqpConnection(mockAuthenticatedPrincipal("user2"),
+                "futureFailure", new Date(System.currentTimeMillis() - 
60000L), failedFuture);
+
+        vhost.registerConnection(thrownFailure);
+        vhost.registerConnection(futureFailure);
+
+        final Map<String, Object> result = vhost.closeIdleConnections(0L, 
true, 2000L);
+
+        assertEquals(2, result.get("matchedCount"), "Unexpected matched 
count");
+        assertEquals(1, result.get("closeRequestedCount"), "Unexpected close 
requested count");
+        assertEquals(0, result.get("closedCount"), "Unexpected closed count");
+        assertEquals(2, result.get("failedCount"), "Unexpected failed count");
+        assertEquals(0, result.get("timedOutCount"), "Unexpected timed out 
count");
+
+        final Map<String, Object> thrownResult = getConnectionResult(result, 
"thrownFailure");
+        assertEquals("FAILED", thrownResult.get("status"), "Unexpected thrown 
failure status");
+        assertEquals("java.lang.IllegalStateException: request failed", 
thrownResult.get("failure"),
+                "Unexpected thrown failure message");
+
+        final Map<String, Object> futureResult = getConnectionResult(result, 
"futureFailure");
+        assertEquals("FAILED", futureResult.get("status"), "Unexpected future 
failure status");
+        assertEquals("java.lang.IllegalStateException: future failed", 
futureResult.get("failure"),
+                "Unexpected future failure message");
+    }
+
+    @Test
+    public void testCloseIdleConnectionsValidatesParameters()
+    {
+        final QueueManagingVirtualHost<?> vhost = 
createVirtualHost(getTestName());
+
+        assertThrows(IllegalArgumentException.class, () ->  
vhost.closeIdleConnections(-1L, false, 2000L),
+                "Exception not thrown for negative idleTimeMillis");
+        assertThrows(IllegalArgumentException.class, () -> 
vhost.closeIdleConnections(0L, false, 0L),
+                "Exception not thrown for non-positive timeoutMillis");
+    }
+
     @Test
     public void testStopVirtualhostClosesConnections()
     {
@@ -636,8 +769,21 @@ public class VirtualHostTest extends UnitTestBase
     }
 
     private AMQPConnection<?> mockAmqpConnection(final Principal principal)
+    {
+        return mockAmqpConnection(principal, getTestName(), new Date(), 
CompletableFuture.completedFuture(null));
+    }
+
+    private AMQPConnection<?> mockAmqpConnection(final Principal principal,
+                                                 final String name,
+                                                 final Date lastMessageTime,
+                                                 final CompletableFuture<Void> 
deleteFuture)
     {
         final AMQPConnection<?> connection = mock(AMQPConnection.class);
+        final String principalName = principal.getName();
+        when(connection.getId()).thenReturn(UUID.randomUUID());
+        when(connection.getName()).thenReturn(name);
+        when(connection.getRemoteAddress()).thenReturn("/127.0.0.1:5672");
+        when(connection.getPrincipal()).thenReturn(principalName);
         when(connection.getAuthorizedPrincipal()).thenReturn(principal);
         final Subject subject =
                 new Subject(true, Set.of(principal), Set.of(), Set.of());
@@ -645,9 +791,27 @@ public class VirtualHostTest extends UnitTestBase
         final CompletableFuture<Void> completableFuture = 
CompletableFuture.completedFuture(null);
         when(connection.closeAsync()).thenReturn(completableFuture);
         when(connection.getCreatedTime()).thenReturn(new Date());
+        when(connection.getLastMessageTime()).thenReturn(lastMessageTime);
+        when(connection.deleteAsync()).thenReturn(deleteFuture);
         return connection;
     }
 
+    private Map<String, Object> getSingleConnectionResult(final Map<String, 
Object> result)
+    {
+        final List<Map<String, Object>> connections = (List<Map<String, 
Object>>) result.get("connections");
+        assertEquals(1, connections.size(), "Unexpected number of connection 
results");
+        return connections.get(0);
+    }
+
+    private Map<String, Object> getConnectionResult(final Map<String, Object> 
result, final String name)
+    {
+        final List<Map<String, Object>> connections = (List<Map<String, 
Object>>) result.get("connections");
+        return connections.stream()
+                .filter(connection -> name.equals(connection.get("name")))
+                .findFirst()
+                .orElseThrow(() -> new AssertionError("Connection result not 
found: " + name));
+    }
+
     private Principal mockAuthenticatedPrincipal(final String principalName)
     {
         final Principal principal = mock(Principal.class);
diff --git 
a/broker-plugins/management-http/src/test/java/org/apache/qpid/server/management/plugin/controller/latest/LatestManagementControllerTest.java
 
b/broker-plugins/management-http/src/test/java/org/apache/qpid/server/management/plugin/controller/latest/LatestManagementControllerTest.java
index 81ce366a4a..bfcad77951 100644
--- 
a/broker-plugins/management-http/src/test/java/org/apache/qpid/server/management/plugin/controller/latest/LatestManagementControllerTest.java
+++ 
b/broker-plugins/management-http/src/test/java/org/apache/qpid/server/management/plugin/controller/latest/LatestManagementControllerTest.java
@@ -27,11 +27,13 @@ import static org.hamcrest.CoreMatchers.is;
 import static org.hamcrest.CoreMatchers.notNullValue;
 import static org.hamcrest.CoreMatchers.nullValue;
 import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.fail;
 import static org.mockito.Mockito.doReturn;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
 
+import java.security.Principal;
 import java.util.Collection;
 import java.util.Date;
 import java.util.HashMap;
@@ -40,6 +42,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.Set;
 import java.util.UUID;
+import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
 import java.util.stream.Collectors;
@@ -55,6 +58,7 @@ import 
org.apache.qpid.server.management.plugin.ManagementException;
 import org.apache.qpid.server.management.plugin.ManagementRequest;
 import org.apache.qpid.server.management.plugin.ManagementResponse;
 import org.apache.qpid.server.management.plugin.RequestType;
+import org.apache.qpid.server.management.plugin.ResponseType;
 import org.apache.qpid.server.model.AuthenticationProvider;
 import org.apache.qpid.server.model.Broker;
 import org.apache.qpid.server.model.BrokerModel;
@@ -69,6 +73,7 @@ import 
org.apache.qpid.server.model.preferences.PreferenceImpl;
 import org.apache.qpid.server.security.SubjectExecutionContext;
 import org.apache.qpid.server.security.auth.AuthenticatedPrincipal;
 import org.apache.qpid.server.security.auth.UsernamePrincipal;
+import org.apache.qpid.server.transport.AMQPConnection;
 import org.apache.qpid.server.virtualhost.QueueManagingVirtualHost;
 import org.apache.qpid.test.utils.UnitTestBase;
 
@@ -460,6 +465,52 @@ public class LatestManagementControllerTest extends 
UnitTestBase
         assertThat(((Number) body).intValue(), is(equalTo(1)));
     }
 
+    @Test
+    public void invokeCloseIdleConnections() throws Exception
+    {
+        final String hostName = "test";
+        final QueueManagingVirtualHost<?> virtualHost = 
createVirtualHostWithQueue(hostName);
+        final AMQPConnection<?> connection = mockAmqpConnection("connection1",
+                new Date(System.currentTimeMillis() - 60000L), 
CompletableFuture.completedFuture(null));
+        virtualHost.registerConnection(connection);
+
+        final List<String> path = List.of(virtualHost.getParent().getName(), 
hostName);
+        final ManagementResponse response = 
_controller.invoke(virtualHost.getBroker(), "virtualhost", path,
+                "closeIdleConnections", Map.of("idleTimeMillis", "0", 
"waitForClose", "true", "timeoutMillis", "2000"),
+                true, true);
+
+        assertThat(response, is(notNullValue()));
+        assertThat(response.getResponseCode(), is(equalTo(200)));
+        assertThat(response.getType(), is(equalTo(ResponseType.DATA)));
+        assertThat(response.getBody(), is(instanceOf(Map.class)));
+
+        final Map<?, ?> result = (Map<?, ?>) response.getBody();
+        assertThat(result.get("matchedCount"), is(equalTo(1)));
+        assertThat(result.get("closeRequestedCount"), is(equalTo(1)));
+        assertThat(result.get("closedCount"), is(equalTo(1)));
+
+        final List<?> connections = (List<?>) result.get("connections");
+        assertThat(connections.size(), is(equalTo(1)));
+        assertThat(connections.get(0), is(instanceOf(Map.class)));
+        assertThat(((Map<?, ?>) connections.get(0)).get("status"), 
is(equalTo("CLOSED")));
+    }
+
+    @Test
+    public void invokeCloseIdleConnectionsUsingGetIsRejected() throws Exception
+    {
+        final String hostName = "test";
+        final QueueManagingVirtualHost<?> virtualHost = 
createVirtualHostWithQueue(hostName);
+        final List<String> path = List.of(virtualHost.getParent().getName(), 
hostName);
+
+        final ManagementException exception = 
assertThrows(ManagementException.class, () ->
+                _controller.invoke(virtualHost.getBroker(), "virtualhost", 
path, "closeIdleConnections",
+                Map.of("idleTimeMillis", "0"), false, true),
+                "Exception not thrown");
+
+        assertThat(exception.getStatusCode(), is(equalTo(405)));
+        assertThat(exception.getHeaders().get("Allow"), is(equalTo("POST")));
+    }
+
     @Test
     public void getPreferences() throws Exception
     {
@@ -721,6 +772,26 @@ public class LatestManagementControllerTest extends 
UnitTestBase
         return virtualHost;
     }
 
+    private AMQPConnection<?> mockAmqpConnection(final String name,
+                                                 final Date lastMessageTime,
+                                                 final CompletableFuture<Void> 
deleteFuture)
+    {
+        final AMQPConnection<?> connection = mock(AMQPConnection.class);
+        final AuthenticationProvider<?> authenticationProvider = 
mock(AuthenticationProvider.class);
+        when(authenticationProvider.getType()).thenReturn("type");
+        when(authenticationProvider.getName()).thenReturn("name");
+
+        final Principal principal = new AuthenticatedPrincipal(new 
UsernamePrincipal("user", authenticationProvider));
+        when(connection.getId()).thenReturn(UUID.randomUUID());
+        when(connection.getName()).thenReturn(name);
+        when(connection.getRemoteAddress()).thenReturn("/127.0.0.1:5672");
+        when(connection.getPrincipal()).thenReturn(principal.getName());
+        when(connection.getAuthorizedPrincipal()).thenReturn(principal);
+        when(connection.getSubject()).thenReturn(new Subject(true, 
Set.of(principal), Set.of(), Set.of()));
+        when(connection.getLastMessageTime()).thenReturn(lastMessageTime);
+        when(connection.deleteAsync()).thenReturn(deleteFuture);
+        return connection;
+    }
 
     private UUID createPreferences(final Subject testSubject,
                                    final QueueManagingVirtualHost<?> 
virtualHost,
diff --git 
a/systests/qpid-systests-http-management/src/test/java/org/apache/qpid/tests/http/rest/model/OperationTest.java
 
b/systests/qpid-systests-http-management/src/test/java/org/apache/qpid/tests/http/rest/model/OperationTest.java
index 5e17fb0902..dc7245af76 100644
--- 
a/systests/qpid-systests-http-management/src/test/java/org/apache/qpid/tests/http/rest/model/OperationTest.java
+++ 
b/systests/qpid-systests-http-management/src/test/java/org/apache/qpid/tests/http/rest/model/OperationTest.java
@@ -36,6 +36,10 @@ import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
 
+import javax.jms.Connection;
+import javax.jms.JMSException;
+import javax.jms.Session;
+
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
@@ -90,6 +94,62 @@ public class OperationTest extends HttpTestBase
         getHelper().submitRequest("virtualhost/getStatistics", "POST", params, 
SC_UNPROCESSABLE_ENTITY);
     }
 
+    @Test
+    public void closeIdleConnectionsWithNoMatchesKeepsConnectionUsable() 
throws Exception
+    {
+        final Connection connection = getConnection();
+        try
+        {
+            final Map<String, Object> result = 
getHelper().postJson("virtualhost/closeIdleConnections",
+                    Map.of("idleTimeMillis", 600000L), MAP_TYPE_REF, SC_OK);
+
+            assertThat(result.get("matchedCount"), is(equalTo(0)));
+            assertThat(result.get("closeRequestedCount"), is(equalTo(0)));
+
+            final Session session = connection.createSession(false, 
Session.AUTO_ACKNOWLEDGE);
+            session.close();
+        }
+        finally
+        {
+            closeQuietly(connection);
+        }
+    }
+
+    @Test
+    public void closeIdleConnectionsWaitsForConnectionClose() throws Exception
+    {
+        final Connection connection = getConnection();
+        try
+        {
+            Thread.sleep(10L);
+
+            final Map<String, Object> result = 
getHelper().postJson("virtualhost/closeIdleConnections",
+                    Map.of("idleTimeMillis", 0L, "waitForClose", true, 
"timeoutMillis", 5000L), MAP_TYPE_REF, SC_OK);
+
+            assertThat(result.get("matchedCount"), is(equalTo(1)));
+            assertThat(result.get("closedCount"), is(equalTo(1)));
+
+            final List<?> connections = (List<?>) result.get("connections");
+            assertThat(connections.size(), is(equalTo(1)));
+            assertThat(((Map<?, ?>) connections.get(0)).get("status"), 
is(equalTo("CLOSED")));
+
+            final Map<String, Object> statistics = 
getHelper().postJson("virtualhost/getStatistics",
+                    Map.of("statistics", List.of("connectionCount")), 
MAP_TYPE_REF, SC_OK);
+            assertThat(((Number) 
statistics.get("connectionCount")).intValue(), is(equalTo(0)));
+        }
+        finally
+        {
+            closeQuietly(connection);
+        }
+    }
+
+    @Test
+    public void closeIdleConnectionsRejectsNegativeIdleTime() throws Exception
+    {
+        getHelper().submitRequest("virtualhost/closeIdleConnections", "POST", 
Map.of("idleTimeMillis", -1L),
+                SC_UNPROCESSABLE_ENTITY);
+    }
+
     @Test
     public void operationNotFound() throws Exception
     {
@@ -140,4 +200,16 @@ public class OperationTest extends HttpTestBase
         assertThat(response, is(notNullValue()));
         assertThat(new String(response, UTF_8).contains("Full thread dump 
captured"), is(equalTo(true)));
     }
+
+    private void closeQuietly(final Connection connection)
+    {
+        try
+        {
+            connection.close();
+        }
+        catch (JMSException e)
+        {
+            // The broker may already have closed the connection under test.
+        }
+    }
 }


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


Reply via email to