chibenwa commented on code in PR #3194:
URL: https://github.com/apache/james-project/pull/3194#discussion_r4090615526
##########
server/queue/queue-activemq/src/main/java/org/apache/james/queue/activemq/EmbeddedActiveMQ.java:
##########
@@ -23,107 +23,87 @@
import jakarta.inject.Inject;
import jakarta.jms.ConnectionFactory;
-import org.apache.activemq.ActiveMQConnectionFactory;
-import org.apache.activemq.ActiveMQPrefetchPolicy;
-import org.apache.activemq.blob.BlobTransferPolicy;
-import org.apache.activemq.broker.BrokerPlugin;
-import org.apache.activemq.broker.BrokerService;
-import org.apache.activemq.broker.jmx.ManagementContext;
-import org.apache.activemq.plugin.StatisticsBrokerPlugin;
-import org.apache.activemq.store.PersistenceAdapter;
+import org.apache.activemq.artemis.api.core.TransportConfiguration;
+import org.apache.activemq.artemis.core.config.Configuration;
+import org.apache.activemq.artemis.core.config.impl.ConfigurationImpl;
+import org.apache.activemq.artemis.core.remoting.impl.invm.InVMAcceptorFactory;
+import org.apache.activemq.artemis.jms.client.ActiveMQConnectionFactory;
import org.apache.james.filesystem.api.FileSystem;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+/**
+ * Embedded Artemis broker replacing the legacy ActiveMQ embedded broker.
+ * Uses Apache ActiveMQ Artemis (Jakarta JMS) as the underlying message broker.
+ */
public class EmbeddedActiveMQ {
private static final Logger LOGGER =
LoggerFactory.getLogger(EmbeddedActiveMQ.class);
- private static final String KAHADB_STORE_LOCATION =
"file://var/store/activemq/brokers/KahaDB";
- private static final String BLOB_TRANSFER_LOCATION =
"file://var/store/activemq/blob-transfer";
- private static final String BROCKERS_LOCATION =
"file://var/store/activemq/brokers";
- private static final String BROKER_ID = "broker";
+ private static final String DATA_DIRECTORY_RELATIVE = "var/store/artemis";
private static final String BROKER_NAME = "james";
- private static final String BROCKER_URI = "tcp://localhost:0";
- private static final String STORE_USAGE_LIMIT_PROPERTY =
"james.activemq.store.usage.limit.bytes";
- private static final String TEMP_USAGE_LIMIT_PROPERTY =
"james.activemq.temp.usage.limit.bytes";
- private static final long DEFAULT_STORE_USAGE_LIMIT_BYTES = 10L * 1024 *
1024 * 1024; // 10 GB
- private static final long DEFAULT_TEMP_USAGE_LIMIT_BYTES = 5L * 1024 *
1024 * 1024; // 5 GB
- private final ActiveMQConnectionFactory activeMQConnectionFactory;
- private final PersistenceAdapter persistenceAdapter;
- private BrokerService brokerService;
+ private final ActiveMQConnectionFactory connectionFactory;
+ private final
org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ
embeddedServer;
@Inject
- private EmbeddedActiveMQ(FileSystem fileSystem, PersistenceAdapter
persistenceAdapter, ActiveMQConfiguration configuration) {
- this.persistenceAdapter = persistenceAdapter;
+ public EmbeddedActiveMQ(FileSystem fileSystem, ActiveMQConfiguration
configuration) {
try {
-
persistenceAdapter.setDirectory(fileSystem.getFile(KAHADB_STORE_LOCATION));
- launchEmbeddedBroker(fileSystem, configuration);
+ String dataDirectory = fileSystem.getFile("file://" +
DATA_DIRECTORY_RELATIVE).getAbsolutePath();
+ embeddedServer = createAndStartBroker(dataDirectory,
configuration);
+ connectionFactory = createConnectionFactory(configuration);
} catch (Exception e) {
- throw new RuntimeException(e);
+ throw new RuntimeException("Failed to start embedded Artemis
broker", e);
}
- activeMQConnectionFactory =
createActiveMQConnectionFactory(createBlobTransferPolicy(fileSystem));
}
public ConnectionFactory getConnectionFactory() {
- return activeMQConnectionFactory;
+ return connectionFactory;
}
@PreDestroy
public void stop() throws Exception {
- LOGGER.info("Stopping embedded ActiveMQ...");
- brokerService.stop();
- LOGGER.info("Stopped embedded ActiveMQ");
+ LOGGER.info("Stopping embedded Artemis broker...");
+ embeddedServer.stop();
+ connectionFactory.close();
+ LOGGER.info("Stopped embedded Artemis broker");
}
- private ActiveMQConnectionFactory
createActiveMQConnectionFactory(BlobTransferPolicy blobTransferPolicy) {
- ActiveMQConnectionFactory connectionFactory = new
ActiveMQConnectionFactory("vm://james?create=false");
- connectionFactory.setTrustAllPackages(false);
- connectionFactory.setBlobTransferPolicy(blobTransferPolicy);
- connectionFactory.setPrefetchPolicy(createActiveMQPrefetchPolicy());
- return connectionFactory;
- }
-
- private ActiveMQPrefetchPolicy createActiveMQPrefetchPolicy() {
- ActiveMQPrefetchPolicy prefetchPolicy = new ActiveMQPrefetchPolicy();
- prefetchPolicy.setQueuePrefetch(0);
- prefetchPolicy.setTopicPrefetch(0);
- return prefetchPolicy;
- }
+ private org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ
createAndStartBroker(String dataDirectory,
+
ActiveMQConfiguration configuration) throws Exception {
+ Configuration config = new ConfigurationImpl()
+ .setSecurityEnabled(false)
+ .setJMXManagementEnabled(false)
+ .setPersistenceEnabled(true)
+ .setJournalDirectory(dataDirectory + "/journal")
+ .setBindingsDirectory(dataDirectory + "/bindings")
+ .setLargeMessagesDirectory(dataDirectory + "/largemessages")
+ .setPagingDirectory(dataDirectory + "/paging")
+
.setJournalSyncTransactional(configuration.isJournalSyncTransactional())
+
.setJournalSyncNonTransactional(configuration.isJournalSyncNonTransactional())
+
.setJournalBufferTimeout_NIO(configuration.getJournalBufferTimeoutNIO())
+ .setJournalBufferSize_NIO(configuration.getJournalBufferSizeNIO())
+ .setJournalMaxIO_NIO(configuration.getJournalMaxIONIO())
+
.setEnabledAsyncConnectionExecution(configuration.isAsyncConnectionExecution())
+ .addAcceptorConfiguration(new
TransportConfiguration(InVMAcceptorFactory.class.getName()))
+ .setName(BROKER_NAME);
- private BlobTransferPolicy createBlobTransferPolicy(FileSystem fileSystem)
{
- FileSystemBlobTransferPolicy blobTransferPolicy = new
FileSystemBlobTransferPolicy();
- blobTransferPolicy.setDefaultUploadUrl(BLOB_TRANSFER_LOCATION);
- blobTransferPolicy.setFileSystem(fileSystem);
- return blobTransferPolicy;
+ org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ
server = new
org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ();
Review Comment:
Use an import
##########
server/queue/queue-activemq/src/main/java/org/apache/james/queue/activemq/EmbeddedActiveMQ.java:
##########
@@ -23,107 +23,87 @@
import jakarta.inject.Inject;
import jakarta.jms.ConnectionFactory;
-import org.apache.activemq.ActiveMQConnectionFactory;
-import org.apache.activemq.ActiveMQPrefetchPolicy;
-import org.apache.activemq.blob.BlobTransferPolicy;
-import org.apache.activemq.broker.BrokerPlugin;
-import org.apache.activemq.broker.BrokerService;
-import org.apache.activemq.broker.jmx.ManagementContext;
-import org.apache.activemq.plugin.StatisticsBrokerPlugin;
-import org.apache.activemq.store.PersistenceAdapter;
+import org.apache.activemq.artemis.api.core.TransportConfiguration;
+import org.apache.activemq.artemis.core.config.Configuration;
+import org.apache.activemq.artemis.core.config.impl.ConfigurationImpl;
+import org.apache.activemq.artemis.core.remoting.impl.invm.InVMAcceptorFactory;
+import org.apache.activemq.artemis.jms.client.ActiveMQConnectionFactory;
import org.apache.james.filesystem.api.FileSystem;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+/**
+ * Embedded Artemis broker replacing the legacy ActiveMQ embedded broker.
+ * Uses Apache ActiveMQ Artemis (Jakarta JMS) as the underlying message broker.
+ */
public class EmbeddedActiveMQ {
private static final Logger LOGGER =
LoggerFactory.getLogger(EmbeddedActiveMQ.class);
- private static final String KAHADB_STORE_LOCATION =
"file://var/store/activemq/brokers/KahaDB";
- private static final String BLOB_TRANSFER_LOCATION =
"file://var/store/activemq/blob-transfer";
- private static final String BROCKERS_LOCATION =
"file://var/store/activemq/brokers";
- private static final String BROKER_ID = "broker";
+ private static final String DATA_DIRECTORY_RELATIVE = "var/store/artemis";
private static final String BROKER_NAME = "james";
- private static final String BROCKER_URI = "tcp://localhost:0";
- private static final String STORE_USAGE_LIMIT_PROPERTY =
"james.activemq.store.usage.limit.bytes";
- private static final String TEMP_USAGE_LIMIT_PROPERTY =
"james.activemq.temp.usage.limit.bytes";
- private static final long DEFAULT_STORE_USAGE_LIMIT_BYTES = 10L * 1024 *
1024 * 1024; // 10 GB
- private static final long DEFAULT_TEMP_USAGE_LIMIT_BYTES = 5L * 1024 *
1024 * 1024; // 5 GB
- private final ActiveMQConnectionFactory activeMQConnectionFactory;
- private final PersistenceAdapter persistenceAdapter;
- private BrokerService brokerService;
+ private final ActiveMQConnectionFactory connectionFactory;
+ private final
org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ
embeddedServer;
@Inject
- private EmbeddedActiveMQ(FileSystem fileSystem, PersistenceAdapter
persistenceAdapter, ActiveMQConfiguration configuration) {
- this.persistenceAdapter = persistenceAdapter;
+ public EmbeddedActiveMQ(FileSystem fileSystem, ActiveMQConfiguration
configuration) {
try {
-
persistenceAdapter.setDirectory(fileSystem.getFile(KAHADB_STORE_LOCATION));
- launchEmbeddedBroker(fileSystem, configuration);
+ String dataDirectory = fileSystem.getFile("file://" +
DATA_DIRECTORY_RELATIVE).getAbsolutePath();
+ embeddedServer = createAndStartBroker(dataDirectory,
configuration);
+ connectionFactory = createConnectionFactory(configuration);
} catch (Exception e) {
- throw new RuntimeException(e);
+ throw new RuntimeException("Failed to start embedded Artemis
broker", e);
}
- activeMQConnectionFactory =
createActiveMQConnectionFactory(createBlobTransferPolicy(fileSystem));
}
public ConnectionFactory getConnectionFactory() {
- return activeMQConnectionFactory;
+ return connectionFactory;
}
@PreDestroy
public void stop() throws Exception {
- LOGGER.info("Stopping embedded ActiveMQ...");
- brokerService.stop();
- LOGGER.info("Stopped embedded ActiveMQ");
+ LOGGER.info("Stopping embedded Artemis broker...");
+ embeddedServer.stop();
+ connectionFactory.close();
+ LOGGER.info("Stopped embedded Artemis broker");
}
- private ActiveMQConnectionFactory
createActiveMQConnectionFactory(BlobTransferPolicy blobTransferPolicy) {
- ActiveMQConnectionFactory connectionFactory = new
ActiveMQConnectionFactory("vm://james?create=false");
- connectionFactory.setTrustAllPackages(false);
- connectionFactory.setBlobTransferPolicy(blobTransferPolicy);
- connectionFactory.setPrefetchPolicy(createActiveMQPrefetchPolicy());
- return connectionFactory;
- }
-
- private ActiveMQPrefetchPolicy createActiveMQPrefetchPolicy() {
- ActiveMQPrefetchPolicy prefetchPolicy = new ActiveMQPrefetchPolicy();
- prefetchPolicy.setQueuePrefetch(0);
- prefetchPolicy.setTopicPrefetch(0);
- return prefetchPolicy;
- }
+ private org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ
createAndStartBroker(String dataDirectory,
+
ActiveMQConfiguration configuration) throws Exception {
+ Configuration config = new ConfigurationImpl()
+ .setSecurityEnabled(false)
+ .setJMXManagementEnabled(false)
+ .setPersistenceEnabled(true)
+ .setJournalDirectory(dataDirectory + "/journal")
+ .setBindingsDirectory(dataDirectory + "/bindings")
+ .setLargeMessagesDirectory(dataDirectory + "/largemessages")
+ .setPagingDirectory(dataDirectory + "/paging")
+
.setJournalSyncTransactional(configuration.isJournalSyncTransactional())
+
.setJournalSyncNonTransactional(configuration.isJournalSyncNonTransactional())
+
.setJournalBufferTimeout_NIO(configuration.getJournalBufferTimeoutNIO())
+ .setJournalBufferSize_NIO(configuration.getJournalBufferSizeNIO())
+ .setJournalMaxIO_NIO(configuration.getJournalMaxIONIO())
+
.setEnabledAsyncConnectionExecution(configuration.isAsyncConnectionExecution())
+ .addAcceptorConfiguration(new
TransportConfiguration(InVMAcceptorFactory.class.getName()))
+ .setName(BROKER_NAME);
- private BlobTransferPolicy createBlobTransferPolicy(FileSystem fileSystem)
{
- FileSystemBlobTransferPolicy blobTransferPolicy = new
FileSystemBlobTransferPolicy();
- blobTransferPolicy.setDefaultUploadUrl(BLOB_TRANSFER_LOCATION);
- blobTransferPolicy.setFileSystem(fileSystem);
- return blobTransferPolicy;
+ org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ
server = new
org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ();
Review Comment:
Use an import
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]