This is an automated email from the ASF dual-hosted git repository.
jungm pushed a commit to branch tomee-10.x
in repository https://gitbox.apache.org/repos/asf/tomee.git
The following commit(s) were added to refs/heads/tomee-10.x by this push:
new 0a2466445f TOMEE-4712 - activate MDB endpoints when the bean is
started instead of deployed (#2951) (#3014)
0a2466445f is described below
commit 0a2466445fb14b4cf48429ea099a8ca6fea912cd
Author: Markus Jung <[email protected]>
AuthorDate: Thu Oct 1 13:29:12 2026 +0200
TOMEE-4712 - activate MDB endpoints when the bean is started instead of
deployed (#2951) (#3014)
* TOMEE-4712 activate MDB endpoints when the bean is started instead of
deployed
* TOMEE-4712 do not deactivate an endpoint whose activation failed
The activation context was marked as started before the endpoint was
activated and stayed started if the resource adapter rejected the activation.
As the endpoint is activated on start now, the bean is still deployed at that
point, so the undeploy of the failed application deactivated an endpoint which
was never active.
---------
Co-authored-by: Richard Zowalla <[email protected]>
---
.../apache/openejb/core/mdb/BaseMdbContainer.java | 8 +
.../org/apache/openejb/core/mdb/MdbContainer.java | 39 ++--
.../openejb/core/mdb/MdbInstanceManager.java | 36 ++--
.../apache/openejb/core/mdb/MdbPoolContainer.java | 6 +
.../openejb/core/mdb/MdbActivationOrderTest.java | 223 +++++++++++++++++++++
5 files changed, 276 insertions(+), 36 deletions(-)
diff --git
a/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/BaseMdbContainer.java
b/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/BaseMdbContainer.java
index 0162a02205..ccd5810301 100644
---
a/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/BaseMdbContainer.java
+++
b/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/BaseMdbContainer.java
@@ -44,6 +44,14 @@ public interface BaseMdbContainer {
Properties getProperties();
+ static boolean isActiveOnStartup(final BeanContext beanContext) {
+ String activeOnStartupSetting =
beanContext.getActivationProperties().get("MdbActiveOnStartup");
+ if (activeOnStartupSetting == null) {
+ activeOnStartupSetting =
beanContext.getActivationProperties().get("DeliveryActive");
+ }
+ return activeOnStartupSetting == null ||
Boolean.parseBoolean(activeOnStartupSetting);
+ }
+
void deploy(BeanContext beanContext) throws OpenEJBException;
void start(BeanContext info) throws OpenEJBException;
diff --git
a/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbContainer.java
b/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbContainer.java
index bda66db61d..3db1654ab0 100644
---
a/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbContainer.java
+++
b/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbContainer.java
@@ -229,30 +229,13 @@ public class MdbContainer implements RpcContainer,
BaseMdbContainer {
}
- // activate the endpoint
+ // the endpoint is activated in start(), once the other beans of the
application are started
CURRENT.set(beanContext);
try {
final MdbActivationContext activationContext = new
MdbActivationContext(Thread.currentThread().getContextClassLoader(),
beanContext, resourceAdapter, endpointFactory, activationSpec);
activationContexts.put(beanContext, activationContext);
- boolean activeOnStartup = true;
- String activeOnStartupSetting =
beanContext.getActivationProperties().get("MdbActiveOnStartup");
-
- if (activeOnStartupSetting == null) {
- activeOnStartupSetting =
beanContext.getActivationProperties().get("DeliveryActive");
- }
-
- if (activeOnStartupSetting != null) {
- activeOnStartup = Boolean.parseBoolean(activeOnStartupSetting);
- }
-
- if (activeOnStartup) {
- activationContext.start();
- } else {
- logger.info("Not auto-activating endpoint for " +
beanContext.getDeploymentID());
- }
-
String jmxName =
beanContext.getActivationProperties().get("MdbJMXControl");
if (jmxName == null) {
jmxName = "true";
@@ -368,6 +351,22 @@ public class MdbContainer implements RpcContainer,
BaseMdbContainer {
}
public void start(final BeanContext info) throws OpenEJBException {
+ final MdbActivationContext activationContext =
activationContexts.get(info);
+ if (activationContext != null) {
+ if (BaseMdbContainer.isActiveOnStartup(info)) {
+ CURRENT.set(info);
+ try {
+ activationContext.start();
+ } catch (final ResourceException e) {
+ throw new OpenEJBException(e);
+ } finally {
+ CURRENT.remove();
+ }
+ } else {
+ logger.info("Not auto-activating endpoint for " +
info.getDeploymentID());
+ }
+ }
+
final EjbTimerService timerService = info.getEjbTimerService();
if (timerService != null) {
timerService.start();
@@ -688,6 +687,10 @@ public class MdbContainer implements RpcContainer,
BaseMdbContainer {
Thread.currentThread().setContextClassLoader(classLoader);
resourceAdapter.endpointActivation(endpointFactory,
activationSpec);
logger.info("Activated endpoint for " +
beanContext.getDeploymentID());
+ } catch (final ResourceException | RuntimeException e) {
+ // the endpoint is not active, so it must not be deactivated
on undeploy
+ started.set(false);
+ throw e;
} finally {
Thread.currentThread().setContextClassLoader(oldCl);
}
diff --git
a/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbInstanceManager.java
b/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbInstanceManager.java
index c0f1fd56cb..3c181e4a0b 100644
---
a/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbInstanceManager.java
+++
b/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbInstanceManager.java
@@ -239,29 +239,12 @@ public class MdbInstanceManager {
}
}
- // activate the endpoint
+ // the endpoint is activated in start(), once the other beans of the
application are started
try {
final MdbPoolContainer.MdbActivationContext activationContext =
new
MdbPoolContainer.MdbActivationContext(Thread.currentThread().getContextClassLoader(),
beanContext, resourceAdapter, endpointFactory, activationSpec);
activationContexts.put(beanContext, activationContext);
- boolean activeOnStartup = true;
- String activeOnStartupSetting =
beanContext.getActivationProperties().get("MdbActiveOnStartup");
-
- if (activeOnStartupSetting == null) {
- activeOnStartupSetting =
beanContext.getActivationProperties().get("DeliveryActive");
- }
-
- if (activeOnStartupSetting != null) {
- activeOnStartup = Boolean.parseBoolean(activeOnStartupSetting);
- }
-
- if (activeOnStartup) {
- activationContext.start();
- } else {
- logger.info("Not auto-activating endpoint for " +
beanContext.getDeploymentID());
- }
-
String jmxControlName =
beanContext.getActivationProperties().get("MdbJMXControl");
if (jmxControlName == null) {
jmxControlName = "true";
@@ -303,6 +286,23 @@ public class MdbInstanceManager {
data.getPool().start();
}
+ public void start(final BeanContext beanContext) throws OpenEJBException {
+ final MdbPoolContainer.MdbActivationContext activationContext =
activationContexts.get(beanContext);
+ if (activationContext == null) {
+ return;
+ }
+
+ if (BaseMdbContainer.isActiveOnStartup(beanContext)) {
+ try {
+ activationContext.start();
+ } catch (final ResourceException e) {
+ throw new OpenEJBException(e);
+ }
+ } else {
+ logger.info("Not auto-activating endpoint for " +
beanContext.getDeploymentID());
+ }
+ }
+
public void undeploy(final BeanContext beanContext) {
final MdbPoolContainer.MdbActivationContext actContext =
activationContexts.get(beanContext);
if (actContext == null) {
diff --git
a/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbPoolContainer.java
b/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbPoolContainer.java
index 45aeae645e..ef738e5042 100644
---
a/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbPoolContainer.java
+++
b/container/openejb-core/src/main/java/org/apache/openejb/core/mdb/MdbPoolContainer.java
@@ -241,6 +241,8 @@ public class MdbPoolContainer implements RpcContainer,
BaseMdbContainer {
}
public void start(final BeanContext info) throws OpenEJBException {
+ instanceManager.start(info);
+
final EjbTimerService timerService = info.getEjbTimerService();
if (timerService != null) {
timerService.start();
@@ -524,6 +526,10 @@ public class MdbPoolContainer implements RpcContainer,
BaseMdbContainer {
Thread.currentThread().setContextClassLoader(classLoader);
resourceAdapter.endpointActivation(endpointFactory,
activationSpec);
logger.info("Activated endpoint for " +
beanContext.getDeploymentID());
+ } catch (final ResourceException | RuntimeException e) {
+ // the endpoint is not active, so it must not be deactivated
on undeploy
+ started.set(false);
+ throw e;
} finally {
Thread.currentThread().setContextClassLoader(oldCl);
}
diff --git
a/container/openejb-core/src/test/java/org/apache/openejb/core/mdb/MdbActivationOrderTest.java
b/container/openejb-core/src/test/java/org/apache/openejb/core/mdb/MdbActivationOrderTest.java
new file mode 100644
index 0000000000..74822affab
--- /dev/null
+++
b/container/openejb-core/src/test/java/org/apache/openejb/core/mdb/MdbActivationOrderTest.java
@@ -0,0 +1,223 @@
+/*
+ * 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.openejb.core.mdb;
+
+import org.apache.openejb.OpenEJBException;
+import org.apache.openejb.assembler.classic.Assembler;
+import org.apache.openejb.assembler.classic.MdbContainerInfo;
+import org.apache.openejb.assembler.classic.ResourceInfo;
+import org.apache.openejb.assembler.classic.SecurityServiceInfo;
+import org.apache.openejb.assembler.classic.TransactionServiceInfo;
+import org.apache.openejb.config.ConfigurationFactory;
+import org.apache.openejb.config.sys.Container;
+import org.apache.openejb.config.sys.Resource;
+import org.apache.openejb.core.ivm.naming.InitContextFactory;
+import org.apache.openejb.jee.EjbJar;
+import org.apache.openejb.jee.MessageDrivenBean;
+import org.apache.openejb.jee.StatelessBean;
+import org.apache.openejb.loader.SystemInstance;
+import org.junit.After;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.Parameterized;
+
+import jakarta.ejb.EJB;
+import jakarta.ejb.MessageDriven;
+import jakarta.ejb.Stateless;
+import jakarta.resource.ResourceException;
+import jakarta.resource.spi.ActivationSpec;
+import jakarta.resource.spi.BootstrapContext;
+import jakarta.resource.spi.InvalidPropertyException;
+import jakarta.resource.spi.endpoint.MessageEndpoint;
+import jakarta.resource.spi.endpoint.MessageEndpointFactory;
+import javax.naming.Context;
+import javax.transaction.xa.XAResource;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.List;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.fail;
+
+/**
+ * A resource adapter may deliver messages as soon as an endpoint is
activated, for example
+ * messages that were left on a destination when the server stopped. The
session beans an
+ * MDB calls must be started by then.
+ */
+@RunWith(Parameterized.class)
+public class MdbActivationOrderTest {
+
+ @After
+ public void tearDown() {
+ final Assembler assembler =
SystemInstance.get().getComponent(Assembler.class);
+ if (assembler != null) {
+ assembler.destroy();
+ }
+ SystemInstance.reset();
+ PendingMessageResourceAdapter.failActivation = false;
+ PendingMessageResourceAdapter.deactivations = 0;
+ }
+
+ @Parameterized.Parameters(name = "pool={0}")
+ public static Collection<Object[]> pool() {
+ return List.of(new Object[]{false}, new Object[]{true});
+ }
+
+ @Parameterized.Parameter
+ public boolean pool;
+
+ @Test
+ public void pendingMessageIsDeliveredToStartedSessionBeans() throws
Exception {
+ createApplication();
+
+ if (ListenerBean.failure != null) {
+ throw new AssertionError("pending message failed",
ListenerBean.failure);
+ }
+ assertEquals(List.of("Hello pending"), ListenerBean.received);
+ }
+
+ @Test
+ public void failedActivationIsNotDeactivated() throws Exception {
+ PendingMessageResourceAdapter.failActivation = true;
+
+ try {
+ createApplication();
+ fail("the application must not start if the endpoint can't be
activated");
+ } catch (final OpenEJBException expected) {
+ // no-op
+ }
+
+ assertEquals(0, PendingMessageResourceAdapter.deactivations);
+ }
+
+ private void createApplication() throws Exception {
+ System.setProperty(Context.INITIAL_CONTEXT_FACTORY,
InitContextFactory.class.getName());
+
+ final ConfigurationFactory config = new ConfigurationFactory();
+ final Assembler assembler = new Assembler();
+
assembler.createTransactionManager(config.configureService(TransactionServiceInfo.class));
+
assembler.createSecurityService(config.configureService(SecurityServiceInfo.class));
+
+ final Resource resourceAdapter = new Resource("PendingRA");
+
resourceAdapter.setClassName(PendingMessageResourceAdapter.class.getName());
+ assembler.createResource(config.configureService(resourceAdapter,
ResourceInfo.class));
+
+ final Container container = new Container("PendingMdbContainer",
"MESSAGE", null);
+ container.getProperties().setProperty("ResourceAdapter", "PendingRA");
+ container.getProperties().setProperty("MessageListenerInterface",
PendingMessageListener.class.getName());
+ container.getProperties().setProperty("ActivationSpecClass",
PendingMessageActivationSpec.class.getName());
+ container.getProperties().setProperty("Pool", Boolean.toString(pool));
+ assembler.createContainer(config.configureService(container,
MdbContainerInfo.class));
+
+ final EjbJar ejbJar = new EjbJar("activation-order");
+ ejbJar.addEnterpriseBean(new StatelessBean(Greeter.class));
+ ejbJar.addEnterpriseBean(new MessageDrivenBean(ListenerBean.class));
+
+ ListenerBean.received.clear();
+ ListenerBean.failure = null;
+
+ assembler.createApplication(config.configureApplication(ejbJar));
+ }
+
+ @Stateless
+ public static class Greeter {
+ public String greet(final String name) {
+ return "Hello " + name;
+ }
+ }
+
+ @MessageDriven
+ public static class ListenerBean implements PendingMessageListener {
+ static final List<String> received = new ArrayList<>();
+ static Throwable failure;
+
+ @EJB
+ private Greeter greeter;
+
+ @Override
+ public void onMessage(final String message) {
+ try {
+ received.add(greeter.greet(message));
+ } catch (final RuntimeException e) {
+ failure = e;
+ }
+ }
+ }
+
+ public interface PendingMessageListener {
+ void onMessage(String message);
+ }
+
+ public static class PendingMessageResourceAdapter implements
jakarta.resource.spi.ResourceAdapter {
+ static boolean failActivation;
+ static int deactivations;
+
+ @Override
+ public void start(final BootstrapContext bootstrapContext) {
+ }
+
+ @Override
+ public void stop() {
+ }
+
+ @Override
+ public void endpointActivation(final MessageEndpointFactory factory,
final ActivationSpec spec) throws ResourceException {
+ if (failActivation) {
+ throw new ResourceException("activation failed");
+ }
+
+ final MessageEndpoint endpoint = factory.createEndpoint(null);
+ try {
+
endpoint.beforeDelivery(PendingMessageListener.class.getMethod("onMessage",
String.class));
+ ((PendingMessageListener) endpoint).onMessage("pending");
+ endpoint.afterDelivery();
+ } catch (final NoSuchMethodException e) {
+ throw new ResourceException(e);
+ } finally {
+ endpoint.release();
+ }
+ }
+
+ @Override
+ public void endpointDeactivation(final MessageEndpointFactory factory,
final ActivationSpec spec) {
+ deactivations++;
+ }
+
+ @Override
+ public XAResource[] getXAResources(final ActivationSpec[] specs) {
+ return new XAResource[0];
+ }
+ }
+
+ public static class PendingMessageActivationSpec implements ActivationSpec
{
+ private jakarta.resource.spi.ResourceAdapter resourceAdapter;
+
+ @Override
+ public void validate() throws InvalidPropertyException {
+ }
+
+ @Override
+ public jakarta.resource.spi.ResourceAdapter getResourceAdapter() {
+ return resourceAdapter;
+ }
+
+ @Override
+ public void setResourceAdapter(final
jakarta.resource.spi.ResourceAdapter resourceAdapter) {
+ this.resourceAdapter = resourceAdapter;
+ }
+ }
+}