From f4e066069d672b3cf10f65bbec064c5cec9f6a29 Mon Sep 17 00:00:00 2001 From: Timothy Bish Date: Tue, 11 Aug 2026 14:02:43 -0400 Subject: [PATCH] ARTEMIS-6184 Return error for remote link close events When a Qpid JMS or other client attempts to unsubscribe a durable subscription but cannot due to some error such as not having permissions to do so the broker does not return an error condition in its response and so the client cannot signal that the unsubscribe request failed. --- .../proton/ProtonServerSenderContext.java | 18 +- .../integration/amqp/JMSDurableUnsubTest.java | 155 ++++++++++++++++++ 2 files changed, 172 insertions(+), 1 deletion(-) create mode 100644 tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/amqp/JMSDurableUnsubTest.java diff --git a/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/proton/ProtonServerSenderContext.java b/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/proton/ProtonServerSenderContext.java index fb837a3b221..def1d4b163f 100644 --- a/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/proton/ProtonServerSenderContext.java +++ b/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/proton/ProtonServerSenderContext.java @@ -16,6 +16,7 @@ */ package org.apache.activemq.artemis.protocol.amqp.proton; +import java.util.Objects; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -44,6 +45,7 @@ import org.apache.qpid.proton.amqp.messaging.Modified; import org.apache.qpid.proton.amqp.messaging.Outcome; import org.apache.qpid.proton.amqp.transaction.TransactionalState; +import org.apache.qpid.proton.amqp.transport.AmqpError; import org.apache.qpid.proton.amqp.transport.DeliveryState; import org.apache.qpid.proton.amqp.transport.DeliveryState.DeliveryStateType; import org.apache.qpid.proton.amqp.transport.ErrorCondition; @@ -292,10 +294,24 @@ public void close(boolean remoteLinkClose) throws ActiveMQAMQPException { sessionSPI.closeSender(brokerConsumer); } // if this is a link close rather than a connection close or detach, we need to delete - // any durable resources for say pub subs + // any durable resources for say pub subs, if the action fails we set the error condition controller.close(remoteLinkClose); } catch (Exception e) { logger.warn(e.getMessage(), e); + final ErrorCondition error = new ErrorCondition(); + + error.setDescription("Error on link close: " + + (Objects.requireNonNullElse(e.getMessage(), e.getClass().getSimpleName()))); + + if (e instanceof ActiveMQAMQPException amqpEx) { + error.setCondition(amqpEx.getAmqpError()); + } else if (e instanceof ActiveMQSecurityException) { + error.setCondition(AmqpError.UNAUTHORIZED_ACCESS); + } else { + error.setCondition(AmqpError.INTERNAL_ERROR); + } + + sender.setCondition(error); } finally { messageWriter.close(); } diff --git a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/amqp/JMSDurableUnsubTest.java b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/amqp/JMSDurableUnsubTest.java new file mode 100644 index 00000000000..71af7d53487 --- /dev/null +++ b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/amqp/JMSDurableUnsubTest.java @@ -0,0 +1,155 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.activemq.artemis.tests.integration.amqp; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import java.util.HashSet; +import java.util.Set; + +import javax.jms.Connection; +import javax.jms.JMSSecurityException; +import javax.jms.MessageConsumer; +import javax.jms.Session; +import javax.jms.Topic; + +import org.apache.activemq.artemis.api.core.QueueConfiguration; +import org.apache.activemq.artemis.api.core.RoutingType; +import org.apache.activemq.artemis.api.core.SimpleString; +import org.apache.activemq.artemis.core.security.Role; +import org.apache.activemq.artemis.core.server.ActiveMQServer; +import org.apache.activemq.artemis.core.server.impl.AddressInfo; +import org.apache.activemq.artemis.core.settings.HierarchicalRepository; +import org.apache.activemq.artemis.spi.core.security.ActiveMQJAASSecurityManager; +import org.apache.activemq.artemis.utils.Wait; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +public class JMSDurableUnsubTest extends JMSClientTestSupport { + + @Override + protected boolean isSecurityEnabled() { + return true; + } + + @Override + protected void createAddressAndQueues(ActiveMQServer server) throws Exception { + // Default Queue + server.addAddressInfo(new AddressInfo(SimpleString.of(getQueueName()), RoutingType.ANYCAST)); + server.createQueue(QueueConfiguration.of(getQueueName()).setRoutingType(RoutingType.ANYCAST)); + + // Default DLQ + server.addAddressInfo(new AddressInfo(SimpleString.of(getDeadLetterAddress()), RoutingType.ANYCAST)); + server.createQueue(QueueConfiguration.of(getDeadLetterAddress()).setRoutingType(RoutingType.ANYCAST)); + + // Create Topic for durable subscriptions to use + server.addAddressInfo(new AddressInfo(SimpleString.of(getTopicName()), RoutingType.MULTICAST)); + } + + @Override + protected void enableSecurity(ActiveMQServer server, String... securityMatches) { + ActiveMQJAASSecurityManager securityManager = (ActiveMQJAASSecurityManager) server.getSecurityManager(); + + securityManager.getConfiguration().addUser(noprivUser, noprivPass); + securityManager.getConfiguration().addRole(noprivUser, "nothing"); + securityManager.getConfiguration().addUser(browseUser, browsePass); + securityManager.getConfiguration().addRole(browseUser, "browser"); + securityManager.getConfiguration().addUser(guestUser, guestPass); + securityManager.getConfiguration().addRole(guestUser, "guest"); + securityManager.getConfiguration().addUser(fullUser, fullPass); + securityManager.getConfiguration().addRole(fullUser, "full"); + + HierarchicalRepository> securityRepository = server.getSecurityRepository(); + Set value = new HashSet<>(); + value.add(new Role("nothing", false, false, false, false, false, false, false, false, false, false, false, false)); + value.add(new Role("browser", false, false, false, false, false, false, false, true, false, false, false, false)); + value.add(new Role("guest", false, true, false, false, false, false, false, true, false, false, false, false)); + value.add(new Role("full", true, true, true, true, true, true, true, true, true, true, false, false)); + securityRepository.addMatch(getTopicName(), value); + + for (String match : securityMatches) { + securityRepository.addMatch(match, value); + } + + server.getConfiguration().setSecurityEnabled(true); + } + + @Test + @Timeout(20) + public void testUnsubscribeAllowedFromAuthorizedUserNoErrorReturned() throws Throwable { + final String clientId = "test-" + getTestName(); + final String subscriptionName = "test-" + getTestName() + "-sub"; + + // Full privilege user creates a subscription and leaves it + try (Connection connection = createConnection(fullUser, fullPass, clientId)) { + Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); + Topic topic = session.createTopic(getTopicName()); + MessageConsumer consumer = session.createDurableSubscriber(topic, subscriptionName); + + connection.start(); + + consumer.close(); + } + + Wait.assertTrue(() -> server.addressQuery(SimpleString.of(getTopicName())).isExists(), 2000, 50); + Wait.assertTrue(() -> server.bindingQuery(SimpleString.of(getTopicName()), false).getQueueNames().size() == 1, 2000, 50); + + // Then it removes it which should not return an error and bindings should be cleaned up + try (Connection connection = createConnection(fullUser, fullPass, clientId)) { + final Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); + + assertDoesNotThrow(() -> session.unsubscribe(subscriptionName)); + } + + Wait.assertTrue(() -> server.addressQuery(SimpleString.of(getTopicName())).isExists(), 2000, 50); + Wait.assertTrue(() -> server.bindingQuery(SimpleString.of(getTopicName()), false).getQueueNames().size() == 0, 2000, 50); + } + + @Test + @Timeout(20) + public void testUnsubscribeDisallowedFromUnauthorizedUserAndErrorReturned() throws Throwable { + final String clientId = "test-" + getTestName(); + final String subscriptionName = "test-" + getTestName() + "-sub"; + + // Full privilege user creates a subscription and leaves it + try (Connection connection = createConnection(fullUser, fullPass, clientId)) { + final Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); + final Topic topic = session.createTopic(getTopicName()); + + session.createDurableSubscriber(topic, subscriptionName); + + connection.start(); + } + + Wait.assertTrue(() -> server.addressQuery(SimpleString.of(getTopicName())).isExists()); + Wait.assertTrue(() -> server.bindingQuery(SimpleString.of(getTopicName()), false).getQueueNames().size() == 1); + + // Low privilege user tries to a remove that subscription but fails and subscription remains + // the client should get a security exception in the response. + try (Connection connection = createConnection(guestUser, guestPass, clientId)) { + final Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); + + assertThrows(JMSSecurityException.class, () -> session.unsubscribe(subscriptionName), + "Expected JMSException when unsubscribing without deleteDurableQueue permission"); + } + + Wait.assertTrue(() -> server.addressQuery(SimpleString.of(getTopicName())).isExists()); + Wait.assertTrue(() -> server.bindingQuery(SimpleString.of(getTopicName()), false).getQueueNames().size() == 1); + } +}