Updated replication connector to use product based host property

This commit is contained in:
Simon Pamies committed 2013-07-12 17:47:03 +02:00
1 parent ae7ebabe45
commit b608a6351e
3 files changed
+19 -10

No files matched your search

@@ -11,6 +11,7 @@ public class Activator implements BundleActivator {
private static final Logger logger = Logger.getLogger(Activator.class.getName());
private static final String PROPERTY_NAME_EXCHANGE_NAME = "replication.exchangeName";
private static final String PROPERTY_NAME_EXCHANGE_HOST = "replication.exchangeHost";
private ReplicationInstancesManager replicationInstancesManager;
@@ -19,11 +20,15 @@ public class Activator implements BundleActivator {
public void start(BundleContext bundleContext) throws Exception {
defaultContext = bundleContext;
String exchangeName = bundleContext.getProperty(PROPERTY_NAME_EXCHANGE_NAME);
String exchangeHost = bundleContext.getProperty(PROPERTY_NAME_EXCHANGE_HOST);
if (exchangeName == null) {
exchangeName = "sapsailinganalytics";
}
if (exchangeHost == null) {
exchangeHost = "localhost";
}
replicationInstancesManager = new ReplicationInstancesManager();
ReplicationService serverReplicationMasterService = new ReplicationServiceImpl(exchangeName, replicationInstancesManager);
ReplicationService serverReplicationMasterService = new ReplicationServiceImpl(exchangeName, exchangeHost, replicationInstancesManager);
bundleContext.registerService(ReplicationService.class, serverReplicationMasterService, null);
logger.info("Registered replication service "+serverReplicationMasterService+" using exchange name "+exchangeName);
}
@@ -71,6 +71,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
* serialized form. Clients need to know this name to be able to bind their queues to the exchange.
*/
private final String exchangeName;
private final String exchangeHost;
/**
* UUID that identifies this server
@@ -80,7 +81,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
private Replicator replicator;
private Thread replicatorThread;
public ReplicationServiceImpl(String exchangeName, final ReplicationInstancesManager replicationInstancesManager) throws IOException {
public ReplicationServiceImpl(String exchangeName, String exchangeHost, final ReplicationInstancesManager replicationInstancesManager) throws IOException {
this.replicationInstancesManager = replicationInstancesManager;
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>();
racingEventServiceTracker = new ServiceTracker<RacingEventService, RacingEventService>(
@@ -88,6 +89,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
racingEventServiceTracker.open();
localService = null;
this.exchangeName = exchangeName;
this.exchangeHost = exchangeHost;
replicator = null;
serverUUID = UUID.randomUUID();
logger.info("Setting " + serverUUID.toString() + " as unique replication identifier.");
@@ -98,20 +100,21 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
* an OSGi service tracker to discover the {@link RacingEventService}, the service to replicate is "injected" here.
* @param exchangeName the name of the exchange to which replicas can bind
*/
public ReplicationServiceImpl(String exchangeName,
public ReplicationServiceImpl(String exchangeName, String exchangeHost,
final ReplicationInstancesManager replicationInstancesManager, RacingEventService localService) throws IOException {
this.replicationInstancesManager = replicationInstancesManager;
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>(); // XXX why is this a map? there should be only one connection to a master
this.localService = localService;
this.exchangeName = exchangeName;
this.exchangeHost = exchangeHost;
replicator = null;
serverUUID = UUID.randomUUID();
logger.info("Setting " + serverUUID.toString() + " as unique replication identifier.");
}
private Channel createMasterChannel(String exchangeName) throws IOException {
private Channel createMasterChannel(String exchangeName, String exchangeHost) throws IOException {
final ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("localhost"); // ...and use default port
connectionFactory.setHost(exchangeHost); // ...and use default port
Channel result = null;
try {
@@ -144,7 +147,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
addAsListenerToRacingEventService();
synchronized (this) {
if (masterChannel == null) {
masterChannel = createMasterChannel(exchangeName);
masterChannel = createMasterChannel(exchangeName, exchangeHost);
}
}
}