Merge branch 'master' into bug1351

This commit is contained in:
Axel Uhl committed 2013-06-21 13:04:38 +02:00
commit 0213ed1d6b
48 files changed
+1064 -171

No files matched your search

+7 -16
View File
@@ -79,7 +79,7 @@ if [ $# -eq 0 ]; then
echo "-n <package name> Name of the bundle you want to hot deploy. Needs fully qualified name like"
echo " com.sap.sailing.monitoring. Only works if there is a fully built server available."
echo "-l <telnet port> Telnet port the OSGi server is running. Optional but enables fully automatic hot-deploy."
echo "-s <target server> Name of server you want to use as target for install or hot-deploy. This overrides default behaviour."
echo "-s <target server> Name of server you want to use as target for install, hot-deploy or remote-reploy. This overrides default behaviour."
echo "-w <ssh target> Target for remote-deploy. Must comply with the following format: user@server."
echo ""
echo "build: builds the server code using Maven to $PROJECT_HOME (log to $START_DIR/build.log)"
@@ -89,8 +89,8 @@ if [ $# -eq 0 ]; then
echo "hot-deploy: performs hot deployment of named bundle into OSGi server"
echo "Example: $0 -n com.sap.sailing.www -l 14888 hot-deploy"
echo ""
echo "remote-deploy: performs hot deployment of named bundle into OSGi server"
echo "Example: $0 -w trac@sapsailing.com remote-deploy"
echo "remote-deploy: performs hot deployment of the java code to a remote server"
echo "Example: $0 -s dev -w trac@sapsailing.com remote-deploy"
echo ""
echo "Active branch is $active_branch"
echo "Project home is $PROJECT_HOME"
@@ -130,7 +130,7 @@ echo INSTALL goes to $ACDIR
shift $((OPTIND-1))
if [[ $@ == "" ]]; then
echo "You need to specify an action [build|install|all|hot-deploy]"
echo "You need to specify an action [build|install|all|hot-deploy|remote-deploy]"
exit 2
fi
@@ -402,16 +402,7 @@ if [[ "$@" == "install" ]] || [[ "$@" == "all" ]]; then
fi
if [[ "$@" == "remote-deploy" ]]; then
read -s -n1 -p "Which server do you want to deploy ([d]ev,[t]est,[p]rod1,p[r]od2): " answer
case $answer in
"D" | "d") SERVER="dev";;
"T" | "t") SERVER="test";;
"P" | "p") SERVER="prod1";;
"R" | "r") SERVER="prod2";;
*) echo "Aborting..."
exit;;
esac
SERVER=$TARGET_SERVER_NAME
echo "Will deploy server $SERVER"
read -s -n1 -p "Did you want me to start a LOCAL build (without tests) for $SERVERS_HOME/$SERVER before deploying (y/n)? " answer
@@ -469,8 +460,8 @@ if [[ "$@" == "remote-deploy" ]]; then
esac
echo ""
$SSH_CMD "$REMOTE_SERVER/stop"
$SSH_CMD "$REMOTE_SERVER/start"
$SSH_CMD "cd $REMOTE_SERVER && $REMOTE_SERVER/stop"
$SSH_CMD "cd $REMOTE_SERVER && $REMOTE_SERVER/start"
echo "Restarted remote server. Please check."
fi
@@ -109,7 +109,7 @@ public class Util {
return !ts.iterator().hasNext();
}
}
public static class Pair<A, B> implements Serializable {
private static final long serialVersionUID = -7631774746419135931L;
@@ -245,5 +245,5 @@ public class Util {
}
return result;
}
}
@@ -40,4 +40,8 @@ public class EventImpl extends EventBaseImpl implements Event {
public void removeRegatta(Regatta regatta) {
regattas.remove(regatta);
}
public String toString() {
return getId() + " " + getName() + " " + getVenue().getName() + " " + isPublic();
}
}
@@ -364,5 +364,9 @@ public class RegattaImpl extends NamedImpl implements Regatta, RaceColumnListene
}
return false;
}
public String toString() {
return getId() + " " + getName() + " " + getScoringScheme().getType().name();
}
}
@@ -1564,5 +1564,9 @@ public abstract class AbstractSimpleLeaderboardImpl implements Leaderboard, Race
}
return result;
}
public String toString() {
return getName() + " " + (getDefaultCourseArea() != null ? getDefaultCourseArea().getName() : "<No course area defined>") + " " + (getScoringScheme() != null ? getScoringScheme().getType().name() : "<No scoring scheme set>");
}
}
@@ -19,4 +19,8 @@ public class EventLeaderboardImpl extends LeaderboardGroupImpl implements EventL
public Event getEvent() {
return event;
}
public String toString() {
return super.toString() + " || " + getEvent().toString();
}
}
@@ -146,4 +146,8 @@ public class LeaderboardGroupImpl implements LeaderboardGroup {
public boolean isDisplayGroupsInReverseOrder() {
return displayGroupsInReverseOrder;
}
public String toString() {
return getName() + " " + getDescription();
}
}
@@ -0,0 +1,21 @@
package com.sap.sailing.util;
import java.io.BufferedReader;
import java.io.File;
import java.io.FileReader;
public class BuildVersion {
public static String getBuildVersion() {
String version = "Unknown or Development (" + System.getProperty("com.sap.sailing.server.name") + ")";
File versionfile = new File(System.getProperty("jetty.home") + File.separator + "version.txt");
if (versionfile.exists()) {
try {
version = new BufferedReader(new FileReader(versionfile)).readLine();
} catch (Exception ex) {
/* ignore */
}
}
return version;
}
}
@@ -1,15 +0,0 @@
<?xml version="1.0" encoding="UTF-8" standalone="no"?>
<launchConfiguration type="org.eclipse.jdt.junit.launchconfig">
<listAttribute key="org.eclipse.debug.core.MAPPED_RESOURCE_PATHS">
<listEntry value="/com.sap.sailing.gwt.ui.test/src/com/sap/sailing/gwt/ui/test/TestColumnSwapping.java"/>
</listAttribute>
<listAttribute key="org.eclipse.debug.core.MAPPED_RESOURCE_TYPES">
<listEntry value="1"/>
</listAttribute>
<stringAttribute key="org.eclipse.jdt.junit.CONTAINER" value=""/>
<booleanAttribute key="org.eclipse.jdt.junit.KEEPRUNNING_ATTR" value="false"/>
<stringAttribute key="org.eclipse.jdt.junit.TESTNAME" value=""/>
<stringAttribute key="org.eclipse.jdt.junit.TEST_KIND" value="org.eclipse.jdt.junit.loader.junit4"/>
<stringAttribute key="org.eclipse.jdt.launching.MAIN_TYPE" value="com.sap.sailing.gwt.ui.test.TestColumnSwapping"/>
<stringAttribute key="org.eclipse.jdt.launching.PROJECT_ATTR" value="com.sap.sailing.gwt.ui.test"/>
</launchConfiguration>
@@ -17,6 +17,17 @@ h1 {
text-align: center;
}
fieldset {
padding: 6px 6px;
margin: 8px 3px 10px 2px;
border: 1px solid black;
}
fieldset legend {
padding: 0px 4px;
font-weight: bold;
}
.chartLegend {
float: left;
margin-right: 4px;
@@ -50,6 +61,14 @@ h1 {
height: 100%;
}
.global-alert-message {
font-weight: bold;
color: red;
font-size: 120%;
text-align: right;
padding: 0px 50px 0px 0px;
}
.errorLabel {
color: #FF0000;
}
@@ -128,6 +147,10 @@ h1 {
background: none repeat scroll 0 0 #e5e5e5;
}
.gwt-Button {
margin: 3px 2px 3px 5px;
}
.gwt-Button:hover {
color: #fff;
}
@@ -34,9 +34,9 @@ public class AdminConsoleEntryPoint extends AbstractEntryPoint implements Regatt
TabPanel tabPanel = new TabPanel();
tabPanel.ensureDebugId("AdministrationTabs");
tabPanel.setAnimationEnabled(true);
tabPanel.setSize("95%", "95%");
tabPanel.setSize("99%", "95%");
rootPanel.add(tabPanel); //, 10, 10);
regattaDisplayers = new HashSet<RegattaDisplayer>();
SailingEventManagementPanel sailingEventManagementPanel = new SailingEventManagementPanel(sailingService, this, stringMessages);
@@ -6,19 +6,23 @@ import com.google.gwt.event.dom.client.ClickEvent;
import com.google.gwt.event.dom.client.ClickHandler;
import com.google.gwt.user.client.rpc.AsyncCallback;
import com.google.gwt.user.client.ui.Button;
import com.google.gwt.user.client.ui.CaptionPanel;
import com.google.gwt.user.client.ui.FlowPanel;
import com.google.gwt.user.client.ui.Grid;
import com.google.gwt.user.client.ui.IntegerBox;
import com.google.gwt.user.client.ui.HorizontalPanel;
import com.google.gwt.user.client.ui.Label;
import com.google.gwt.user.client.ui.TextBox;
import com.google.gwt.user.client.ui.VerticalPanel;
import com.google.gwt.user.client.ui.Widget;
import com.sap.sailing.domain.common.impl.Util.Pair;
import com.sap.sailing.domain.common.impl.Util.Triple;
import com.sap.sailing.gwt.ui.client.DataEntryDialog;
import com.sap.sailing.gwt.ui.client.DataEntryDialog.DialogCallback;
import com.sap.sailing.gwt.ui.client.ErrorReporter;
import com.sap.sailing.gwt.ui.client.IntegerBox;
import com.sap.sailing.gwt.ui.client.SailingServiceAsync;
import com.sap.sailing.gwt.ui.client.StringMessages;
import com.sap.sailing.gwt.ui.client.shared.panels.UserStatusPanel;
import com.sap.sailing.gwt.ui.shared.ReplicaDTO;
import com.sap.sailing.gwt.ui.shared.ReplicationMasterDTO;
import com.sap.sailing.gwt.ui.shared.ReplicationStateDTO;
@@ -32,17 +36,21 @@ import com.sap.sailing.gwt.ui.shared.ReplicationStateDTO;
*/
public class ReplicationPanel extends FlowPanel {
private final Grid registeredReplicas;
private final Grid registeredMasters;
private final SailingServiceAsync sailingService;
private final ErrorReporter errorReporter;
private final StringMessages stringMessages;
private final Button addButton;
private final Button stopReplicationButton;
private final Button removeAllReplicas;
public ReplicationPanel(SailingServiceAsync sailingService, ErrorReporter errorReporter, StringMessages stringMessages) {
this.sailingService = sailingService;
this.stringMessages = stringMessages;
this.errorReporter = errorReporter;
registeredReplicas = new Grid();
registeredReplicas.resizeColumns(3);
add(registeredReplicas);
Button refreshButton = new Button(stringMessages.refresh());
refreshButton.addClickHandler(new ClickHandler() {
@Override
@@ -51,14 +59,90 @@ public class ReplicationPanel extends FlowPanel {
}
});
add(refreshButton);
Button addButton = new Button(stringMessages.add());
final CaptionPanel mastergroup = new CaptionPanel(stringMessages.explainReplicasRegistered());
final VerticalPanel masterpanel = new VerticalPanel();
registeredReplicas = new Grid();
registeredReplicas.resizeColumns(3);
masterpanel.add(registeredReplicas);
removeAllReplicas = new Button(stringMessages.stopAllReplicas());
removeAllReplicas.addClickHandler(new ClickHandler() {
@Override
public void onClick(ClickEvent event) {
stopAllReplicas();
};
});
removeAllReplicas.setEnabled(false);
masterpanel.add(removeAllReplicas);
mastergroup.add(masterpanel);
add(mastergroup);
final CaptionPanel replicagroup = new CaptionPanel(stringMessages.explainConnectionsToMaster());
final VerticalPanel replicapanel = new VerticalPanel();
final HorizontalPanel replicapanelbuttons = new HorizontalPanel();
registeredMasters = new Grid();
registeredMasters.resizeColumns(3);
replicapanel.add(registeredMasters);
addButton = new Button(stringMessages.connectToMaster());
addButton.addClickHandler(new ClickHandler() {
@Override
public void onClick(ClickEvent event) {
addReplication();
}
});
add(addButton);
replicapanelbuttons.add(addButton);
stopReplicationButton = new Button(stringMessages.stopConnectionToMaster());
stopReplicationButton.addClickHandler(new ClickHandler() {
@Override
public void onClick(ClickEvent event) {
stopReplication();
};
});
stopReplicationButton.setEnabled(false);
replicapanelbuttons.add(stopReplicationButton);
replicapanel.add(replicapanelbuttons);
replicagroup.add(replicapanel);
add(replicagroup);
updateReplicaList();
}
protected void stopAllReplicas() {
sailingService.stopAllReplicas(new AsyncCallback<Void>() {
@Override
public void onFailure(Throwable caught) {
errorReporter.reportError(caught.getMessage());
updateReplicaList();
}
@Override
public void onSuccess(Void result) {
removeAllReplicas.setEnabled(false);
updateReplicaList();
}
});
}
private void stopReplication() {
stopReplicationButton.setEnabled(false);
sailingService.stopReplicatingFromMaster(new AsyncCallback<Void>() {
@Override
public void onFailure(Throwable caught) {
errorReporter.reportError(caught.getMessage());
stopReplicationButton.setEnabled(true);
}
@Override
public void onSuccess(Void result) {
addButton.setEnabled(true);
stopReplicationButton.setEnabled(false);
updateReplicaList();
}
});
}
private void addReplication() {
@@ -66,19 +150,28 @@ public class ReplicationPanel extends FlowPanel {
new DialogCallback<Triple<Pair<String, String>, Integer, Integer>>() {
@Override
public void ok(final Triple<Pair<String, String>, Integer, Integer> masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber) {
registeredMasters.removeRow(0);
registeredMasters.insertRow(0);
registeredMasters.setWidget(0, 0, new Label(stringMessages.loading()));
addButton.setEnabled(false);
stopReplicationButton.setEnabled(false);
sailingService.startReplicatingFromMaster(masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getA(),
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getB(),
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getC(),
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getB(), new AsyncCallback<Void>() {
@Override
public void onFailure(Throwable e) {
addButton.setEnabled(true);
errorReporter.reportError(stringMessages.errorStartingReplication(
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getA(),
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getB(), e.getMessage()));
updateReplicaList();
}
@Override
public void onSuccess(Void arg0) {
addButton.setEnabled(false);
stopReplicationButton.setEnabled(true);
updateReplicaList();
}
});
@@ -96,21 +189,33 @@ public class ReplicationPanel extends FlowPanel {
sailingService.getReplicaInfo(new AsyncCallback<ReplicationStateDTO>() {
@Override
public void onSuccess(ReplicationStateDTO replicas) {
int i=0;
while (registeredReplicas.getRowCount() > 0) {
registeredReplicas.removeRow(0);
}
int i=0;
final ReplicationMasterDTO replicatingFromMaster = replicas.getReplicatingFromMaster();
if (replicatingFromMaster != null) {
boolean replicaRegistered = false;
for (final ReplicaDTO replica : replicas.getReplicas()) {
registeredReplicas.insertRow(i);
registeredReplicas.setWidget(i, 0, new Label(stringMessages.replicatingFromMaster(replicatingFromMaster.getHostname(),
replicatingFromMaster.getMessagingPort(), replicatingFromMaster.getServletPort())));
i++;
}
for (ReplicaDTO replica : replicas.getReplicas()) {
registeredReplicas.insertRow(i);
registeredReplicas.setWidget(i, 0, new Label(replica.getHostname()));
registeredReplicas.setWidget(i, 0, new Label((i+1) + ". " + replica.getHostname() + " (" + replica.getIdentifier() + ")"));
registeredReplicas.setWidget(i, 1, new Label(stringMessages.registeredAt(replica.getRegistrationTime().toString())));
final Button removeReplicaButton = new Button(stringMessages.dropReplicaConnection());
removeReplicaButton.addClickHandler(new ClickHandler() {
@Override
public void onClick(ClickEvent event) {
sailingService.stopSingleReplicaInstance(replica.getIdentifier(), new AsyncCallback<Void>() {
@Override
public void onFailure(Throwable caught) {
errorReporter.reportError(caught.getMessage());
updateReplicaList();
}
@Override
public void onSuccess(Void result) {
updateReplicaList();
}
});
}
});
registeredReplicas.setWidget(i, 2, removeReplicaButton);
i++;
for (Map.Entry<String, Integer> e : replica.getOperationCountByOperationClassName().entrySet()) {
registeredReplicas.insertRow(i);
@@ -118,7 +223,43 @@ public class ReplicationPanel extends FlowPanel {
registeredReplicas.setWidget(i, 2, new Label(e.getValue().toString()));
i++;
}
replicaRegistered = true;
}
if (!replicaRegistered) {
registeredReplicas.insertRow(i);
registeredReplicas.setWidget(i, 0, new Label(stringMessages.explainNoConnectionsFromReplicas()));
removeAllReplicas.setEnabled(false);
} else {
removeAllReplicas.setEnabled(true);
}
while (registeredMasters.getRowCount() > 0) {
registeredMasters.removeRow(0);
}
i = 0;
registeredMasters.insertRow(i);
registeredMasters.setWidget(i, 0, new Label("Client UUID: " + replicas.getServerIdentifier()));
i++;
final ReplicationMasterDTO replicatingFromMaster = replicas.getReplicatingFromMaster();
if (replicatingFromMaster != null) {
UserStatusPanel.setGlobalAlert(stringMessages.warningServerIsReplica());
registeredMasters.insertRow(i);
registeredMasters.setWidget(i, 0, new Label(stringMessages.replicatingFromMaster(replicatingFromMaster.getHostname(),
replicatingFromMaster.getMessagingPort(), replicatingFromMaster.getServletPort())));
i++;
addButton.setEnabled(false);
stopReplicationButton.setEnabled(true);
} else {
UserStatusPanel.setGlobalAlert("");
registeredMasters.insertRow(i);
registeredMasters.setWidget(i, 0, new Label(stringMessages.explainNoConnectionsToMaster()));
addButton.setEnabled(true);
stopReplicationButton.setEnabled(false);
}
}
@Override
@@ -143,7 +284,7 @@ public class ReplicationPanel extends FlowPanel {
public AddReplicationDialog(final Validator<Triple<Pair<String, String>, Integer, Integer>> validator,
final DialogCallback<Triple<Pair<String, String>, Integer, Integer>> callback) {
super(stringMessages.add(), stringMessages.enterMaster(),
super(stringMessages.connect(), stringMessages.enterMaster(),
stringMessages.ok(), stringMessages.cancel(), validator, callback);
hostnameEntryField = createTextBox("localhost");
exchangenameEntryField = createTextBox("sapsailinganalytics");
@@ -157,15 +298,19 @@ public class ReplicationPanel extends FlowPanel {
*/
@Override
protected Widget getAdditionalWidget() {
Grid grid = new Grid(4, 2);
Grid grid = new Grid(8, 2);
grid.setWidget(0, 0, new Label(stringMessages.hostname()));
grid.setWidget(0, 1, hostnameEntryField);
grid.setWidget(1, 0, new Label(stringMessages.exchangeName()));
grid.setWidget(1, 1, exchangenameEntryField);
grid.setWidget(2, 0, new Label(stringMessages.messagingPortNumber()));
grid.setWidget(2, 1, messagingPortField);
grid.setWidget(3, 0, new Label(stringMessages.servletPortNumber()));
grid.setWidget(3, 1, servletPortField);
grid.setWidget(1, 0, new Label(stringMessages.explainReplicationHostname()));
grid.setWidget(2, 0, new Label(stringMessages.exchangeName()));
grid.setWidget(2, 1, exchangenameEntryField);
grid.setWidget(3, 0, new Label(stringMessages.explainReplicationExchangeName()));
grid.setWidget(4, 0, new Label(stringMessages.messagingPortNumber()));
grid.setWidget(4, 1, messagingPortField);
grid.setWidget(5, 0, new Label(stringMessages.explainReplicationExchangePort()));
grid.setWidget(6, 0, new Label(stringMessages.servletPortNumber()));
grid.setWidget(6, 1, servletPortField);
grid.setWidget(7, 0, new Label(stringMessages.explainReplicationServletPort()));
return grid;
}
@@ -14,7 +14,6 @@ import com.google.gwt.event.logical.shared.ValueChangeHandler;
import com.google.gwt.user.client.ui.Button;
import com.google.gwt.user.client.ui.CheckBox;
import com.google.gwt.user.client.ui.Grid;
import com.google.gwt.user.client.ui.IntegerBox;
import com.google.gwt.user.client.ui.Label;
import com.google.gwt.user.client.ui.ListBox;
import com.google.gwt.user.client.ui.TextBox;
@@ -23,6 +22,7 @@ import com.google.gwt.user.client.ui.Widget;
import com.sap.sailing.domain.common.Color;
import com.sap.sailing.domain.common.dto.FleetDTO;
import com.sap.sailing.gwt.ui.client.DataEntryDialog;
import com.sap.sailing.gwt.ui.client.IntegerBox;
import com.sap.sailing.gwt.ui.client.StringMessages;
import com.sap.sailing.gwt.ui.shared.SeriesDTO;
@@ -19,7 +19,6 @@ import com.google.gwt.user.client.ui.DoubleBox;
import com.google.gwt.user.client.ui.FlowPanel;
import com.google.gwt.user.client.ui.HasVerticalAlignment;
import com.google.gwt.user.client.ui.HorizontalPanel;
import com.google.gwt.user.client.ui.IntegerBox;
import com.google.gwt.user.client.ui.Label;
import com.google.gwt.user.client.ui.ListBox;
import com.google.gwt.user.client.ui.LongBox;
@@ -0,0 +1,28 @@
package com.sap.sailing.gwt.ui.client;
import java.io.IOException;
import com.google.gwt.dom.client.Document;
import com.google.gwt.text.client.IntegerParser;
import com.google.gwt.text.shared.Renderer;
import com.google.gwt.user.client.ui.ValueBox;
public class IntegerBox extends ValueBox<Integer> {
private static final Renderer<Integer> RENDERER = new Renderer<Integer>() {
public String render(Integer object) {
if (object == null) {
return null;
}
StringBuilder sb = new StringBuilder(String.valueOf(object));
return sb.toString();
}
public void render(Integer object, Appendable appendable) throws IOException {
appendable.append(render(object));
}
};
public IntegerBox() {
super(Document.get().createTextInputElement(), RENDERER, IntegerParser.instance());
}
}
@@ -309,4 +309,10 @@ public interface SailingService extends RemoteService {
List<String> visibleRegattas, boolean showOnlyCurrentlyRunningRaces, boolean showOnlyRacesOfSameDay);
String getBuildVersion();
void stopReplicatingFromMaster();
void stopAllReplicas();
void stopSingleReplicaInstance(String identifier);
}
@@ -423,12 +423,18 @@ public interface SailingServiceAsync {
void updateRegatta(RegattaIdentifier regattaIdentifier, String defaultCourseAreaId, AsyncCallback<Void> callback);
void getBuildVersion(AsyncCallback<String> callback);
void stopReplicatingFromMaster(AsyncCallback<Void> asyncCallback);
void getRegattaStructureForEvent(String eventIdAsString, AsyncCallback<List<RaceGroupDTO>> asyncCallback);
void getRaceStateEntriesForRaceGroup(String eventIdAsString, List<String> visibleCourseAreas,
List<String> visibleRegattas, boolean showOnlyCurrentlyRunningRaces, boolean showOnlyRacesOfSameDay,
AsyncCallback<List<RegattaOverviewEntryDTO>> markedAsyncCallback);
void stopAllReplicas(AsyncCallback<Void> asyncCallback);
void stopSingleReplicaInstance(String identifier, AsyncCallback<Void> asyncCallback);
}
@@ -656,6 +656,18 @@ public interface StringMessages extends Messages {
String startsWithZeroScore();
String regattaOverviewConfiguration();
String firstRaceIsNonDiscardableCarryForward();
String addReplicationMaster();
String connect();
String connectToMaster();
String explainReplicationHostname();
String explainReplicationExchangeName();
String explainReplicationExchangePort();
String explainReplicationServletPort();
String explainReplicasRegistered();
String explainConnectionsToMaster();
String explainNoConnectionsToMaster();
String explainNoConnectionsFromReplicas();
String stopConnectionToMaster();
String asLink();
String raceIsRunning();
String raceIsRunningWithEarlyStarters();
@@ -670,4 +682,7 @@ public interface StringMessages extends Messages {
String raceAbandoned();
String raceAbandonedNoMoreRacingToday();
String raceAbandonedFurtherSignalsAshore();
String stopAllReplicas();
String warningServerIsReplica();
String dropReplicaConnection();
}
@@ -649,6 +649,18 @@ regattaDefinesResultDiscardingRules=Regatta defines result discarding rules
startsWithZeroScore=Series starts with zero score
regattaOverviewConfiguration=Regatta Overview Configuration
firstRaceIsNonDiscardableCarryForward=Starts with non-discardable carry-forward
addReplicationMaster=Connect to replication master server
connect=Connect
connectToMaster=Connect to Master
explainReplicationHostname=Hostname where the replication master is running on. This is also the host where the messaging queue is running on.
explainReplicationExchangeName=The name of the exchange queue. It must be the same as configured in master.
explainReplicationExchangePort=The port of the messaging exchange queue. This is used for all data after initial load.
explainReplicationServletPort=Servlet port number that is used to connect to register this replica and receive initial load.
explainReplicasRegistered=Replicas registered (this server is master):
explainConnectionsToMaster=Connections to a master (this server is a replica):
explainNoConnectionsToMaster=No connections to a master server configured. Click the Connect to Master button to configure one.
explainNoConnectionsFromReplicas=No connections from a replica found. You can connect replicas by using their AdminConsole.
stopConnectionToMaster=Stop Replication
asLink=As link
raceIsRunning=Race is running
raceIsRunningWithEarlyStarters=Race is running (had early starters)
@@ -662,4 +674,7 @@ startPostponedNoMoreRacingToday=Start postponed - No more racing today
startPostponedFurtherSignalsAshore=Start postponed - Further signals ashore
raceAbandoned=Race abandoned
raceAbandonedNoMoreRacingToday=Race abandoned - No more racing today
raceAbandonedFurtherSignalsAshore=Race abandoned - Further signals ashore
raceAbandonedFurtherSignalsAshore=Race abandoned - Further signals ashore
stopAllReplicas=Drop connection to all replicas (Dangerous!)
warningServerIsReplica=ATTENTION: This server is configured to be a replica. Do NOT change anything or this will lead to strange effects!
dropReplicaConnection=Drop connection
@@ -647,17 +647,20 @@ regattaDefinesResultDiscardingRules=Streichergebnisse werden von der Regatta fes
startsWithZeroScore=Serie beginnt mit auf 0 gesetzten Ergebnissen
regattaOverviewConfiguration=Regattaübersicht-Konfiguration
firstRaceIsNonDiscardableCarryForward=Beginnt mit nicht streichbarem Übertrag
asLink=Als Link
raceIsRunning=Rennen läuft
raceIsRunningWithEarlyStarters=Rennen läuft (es gab Frühstarter)
raceIsFinishing=Rennen am beenden
raceIsFinished=Rennen beendet
raceIsScheduled=Startzeit angesetzt
raceIsInStartphase=Rennen in der Startphase
generalRecall=Gesamtrückruf
startPostponed=Startverschiebung
startPostponedNoMoreRacingToday=Startverschiebung - heute keine Rennen mehr
startPostponedFurtherSignalsAshore=Startverschiebung - weitere Signale an Land
raceAbandoned=Rennabbruch
addReplicationMaster=Verbindung mit Master herstellen
connect=Verbinden
connectToMaster=Mit einem Master verbinden
explainReplicationHostname=Der Name oder die IP des Servers auf dem der Master und die Queue laufen.
explainReplicationExchangeName=Name der Messaging Queue.
explainReplicationExchangePort=Port der Messaging Queue.
explainReplicationServletPort=Servlet Port.
explainReplicasRegistered=Registrierte Replicas (dieser Server ist der Master):
explainConnectionsToMaster=Verbindungen zu einem Master (dieser Server ist eine Replica)
explainNoConnectionsToMaster=Es sind keine Verbindungen zu einem Master Server konfiguriert.
explainNoConnectionsFromReplicas=Es sind keine Verbindugen von Replicas registriert.
stopConnectionToMaster=Replikation anhalten
raceAbandonedNoMoreRacingToday=Rennabbruch - heute keine Rennen mehr
raceAbandonedFurtherSignalsAshore=Rennabbruch - weitere Signale an Land
raceAbandonedFurtherSignalsAshore=Rennabbruch - weitere Signale an Land
stoppAllReplicas=Verbindung zu allen Replicas abbrechen
warningServerIsReplica=ACHTUNG: Dieser Server ist als Replica konfiguriert. Vermeiden Sie alle Operationen!
dropReplicaConnection=Verbindung abbrechen
@@ -12,16 +12,16 @@ import com.google.gwt.user.client.ui.CheckBox;
import com.google.gwt.user.client.ui.FocusWidget;
import com.google.gwt.user.client.ui.Grid;
import com.google.gwt.user.client.ui.HorizontalPanel;
import com.google.gwt.user.client.ui.IntegerBox;
import com.google.gwt.user.client.ui.Label;
import com.google.gwt.user.client.ui.VerticalPanel;
import com.google.gwt.user.client.ui.Widget;
import com.sap.sailing.domain.common.WindSourceType;
import com.sap.sailing.gwt.ui.client.DataEntryDialog;
import com.sap.sailing.gwt.ui.client.DataEntryDialog.Validator;
import com.sap.sailing.gwt.ui.client.shared.components.SettingsDialogComponent;
import com.sap.sailing.gwt.ui.client.IntegerBox;
import com.sap.sailing.gwt.ui.client.StringMessages;
import com.sap.sailing.gwt.ui.client.WindSourceTypeFormatter;
import com.sap.sailing.gwt.ui.client.shared.components.SettingsDialogComponent;
public class WindChartSettingsDialogComponent implements SettingsDialogComponent<WindChartSettings> {
private final WindChartSettings initialSettings;
@@ -1,11 +1,11 @@
package com.sap.sailing.gwt.ui.client.shared.filter;
import com.google.gwt.user.client.ui.HorizontalPanel;
import com.google.gwt.user.client.ui.IntegerBox;
import com.google.gwt.user.client.ui.ListBox;
import com.google.gwt.user.client.ui.Widget;
import com.sap.sailing.domain.common.filter.BinaryOperator;
import com.sap.sailing.gwt.ui.client.DataEntryDialog;
import com.sap.sailing.gwt.ui.client.IntegerBox;
public class CompetitorRaceRankFilterUIFactory extends AbstractCompetitorNumberFilterUIFactory<Integer> {
private IntegerBox valueIntegerBox;
@@ -1,12 +1,12 @@
package com.sap.sailing.gwt.ui.client.shared.filter;
import com.google.gwt.user.client.ui.HorizontalPanel;
import com.google.gwt.user.client.ui.IntegerBox;
import com.google.gwt.user.client.ui.ListBox;
import com.google.gwt.user.client.ui.Widget;
import com.sap.sailing.domain.common.dto.CompetitorDTO;
import com.sap.sailing.domain.common.filter.BinaryOperator;
import com.sap.sailing.gwt.ui.client.DataEntryDialog;
import com.sap.sailing.gwt.ui.client.IntegerBox;
public class CompetitorTotalRankFilterUIFactory extends AbstractCompetitorNumberFilterUIFactory<Integer> {
private IntegerBox valueIntegerBox;
@@ -23,6 +23,8 @@ public class UserStatusPanel extends FlowPanel {
private final Label userRolesText;
private final Button logoutButton;
public static final Label alert = new Label("");
public UserStatusPanel(final UserManagementServiceAsync userManagementService, final ErrorReporter errorReporter) {
super();
@@ -68,7 +70,10 @@ public class UserStatusPanel extends FlowPanel {
addFloatingWidget(userNameLabel);
addFloatingWidget(userNameText);
addFloatingWidget(userRolesText);
add(logoutButton);
//add(logoutButton);
alert.setStyleName("global-alert-message");
add(alert);
updateUser(user);
}
@@ -92,4 +97,8 @@ public class UserStatusPanel extends FlowPanel {
logoutButton.setEnabled(false);
}
}
public static void setGlobalAlert(String message) {
alert.setText(message);
}
}
@@ -15,7 +15,6 @@ import com.google.gwt.user.client.ui.FocusWidget;
import com.google.gwt.user.client.ui.Grid;
import com.google.gwt.user.client.ui.HasVerticalAlignment;
import com.google.gwt.user.client.ui.HorizontalPanel;
import com.google.gwt.user.client.ui.IntegerBox;
import com.google.gwt.user.client.ui.Label;
import com.google.gwt.user.client.ui.LongBox;
import com.google.gwt.user.client.ui.RadioButton;
@@ -26,6 +25,7 @@ import com.sap.sailing.domain.common.impl.Util;
import com.sap.sailing.gwt.ui.client.DataEntryDialog;
import com.sap.sailing.gwt.ui.client.DataEntryDialog.Validator;
import com.sap.sailing.gwt.ui.client.DetailTypeFormatter;
import com.sap.sailing.gwt.ui.client.IntegerBox;
import com.sap.sailing.gwt.ui.client.StringMessages;
import com.sap.sailing.gwt.ui.client.shared.components.SettingsDialogComponent;
import com.sap.sailing.gwt.ui.leaderboard.LeaderboardSettings.RaceColumnSelectionStrategies;
@@ -1,8 +1,5 @@
package com.sap.sailing.gwt.ui.server;
import java.io.BufferedReader;
import java.io.File;
import java.io.FileReader;
import java.io.IOException;
import java.io.ObjectOutputStream;
import java.io.Serializable;
@@ -274,10 +271,11 @@ import com.sap.sailing.server.operationaltransformation.UpdateLeaderboardScoreCo
import com.sap.sailing.server.operationaltransformation.UpdateRaceDelayToLive;
import com.sap.sailing.server.operationaltransformation.UpdateSeries;
import com.sap.sailing.server.operationaltransformation.UpdateSpecificRegatta;
import com.sap.sailing.server.replication.ReplicaDescriptor;
import com.sap.sailing.server.replication.ReplicationFactory;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
import com.sap.sailing.server.replication.ReplicationService;
import com.sap.sailing.server.replication.impl.ReplicaDescriptor;
import com.sap.sailing.util.BuildVersion;
/**
* The server side implementation of the RPC service.
@@ -2433,7 +2431,7 @@ public class SailingServiceImpl extends ProxiedRemoteServiceServlet implements S
for (Map.Entry<Class<? extends RacingEventServiceOperation<?>>, Integer> e : statistics.entrySet()) {
replicationCountByOperationClassName.put(e.getKey().getName(), e.getValue());
}
replicaDTOs.add(new ReplicaDTO(replicaDescriptor.getIpAddress().getHostName(), replicaDescriptor.getRegistrationTime().asDate(),
replicaDTOs.add(new ReplicaDTO(replicaDescriptor.getIpAddress().getHostName(), replicaDescriptor.getRegistrationTime().asDate(), replicaDescriptor.getUuid().toString(),
replicationCountByOperationClassName));
}
ReplicationMasterDTO master;
@@ -2444,13 +2442,16 @@ public class SailingServiceImpl extends ProxiedRemoteServiceServlet implements S
master = new ReplicationMasterDTO(replicatingFromMaster.getHostname(), replicatingFromMaster.getMessagingPort(),
replicatingFromMaster.getServletPort());
}
return new ReplicationStateDTO(master, replicaDTOs);
return new ReplicationStateDTO(master, replicaDTOs, service.getServerIdentifier().toString());
}
@Override
public void startReplicatingFromMaster(String masterName, String exchangeName, int servletPort, int messagingPort) throws IOException, ClassNotFoundException, InterruptedException {
// the queue name must be always the same for this server. in order to achieve
// this we're using the unique server identifier
getReplicationService().startToReplicateFrom(
ReplicationFactory.INSTANCE.createReplicationMasterDescriptor(masterName, exchangeName, servletPort, messagingPort));
ReplicationFactory.INSTANCE.createReplicationMasterDescriptor(masterName, exchangeName, servletPort, messagingPort,
getReplicationService().getServerIdentifier().toString()));
}
@Override
@@ -3078,16 +3079,40 @@ public class SailingServiceImpl extends ProxiedRemoteServiceServlet implements S
return entry;
}
@Override
public String getBuildVersion() {
String version = "Unknown or Development";
File versionfile = new File(System.getProperty("jetty.home") + File.separator + "version.txt");
if (versionfile.exists()) {
try {
version = new BufferedReader(new FileReader(versionfile)).readLine();
} catch (Exception ex) {
logger.severe("Could not load file " + versionfile.getAbsolutePath());
}
return BuildVersion.getBuildVersion();
}
@Override
public void stopReplicatingFromMaster() {
try {
getReplicationService().stopToReplicateFromMaster();
} catch (IOException e) {
e.printStackTrace();
throw new RuntimeException(e);
}
}
@Override
public void stopAllReplicas() {
try {
getReplicationService().stopAllReplica();
} catch (IOException e) {
e.printStackTrace();
throw new RuntimeException(e);
}
}
@Override
public void stopSingleReplicaInstance(String identifier) {
UUID uuid = UUID.fromString(identifier);
ReplicaDescriptor replicaDescriptor = new ReplicaDescriptor(null, uuid, "");
try {
getReplicationService().unregisterReplica(replicaDescriptor);
} catch (IOException e) {
e.printStackTrace();
throw new RuntimeException(e);
}
return version;
}
}
@@ -7,12 +7,14 @@ import com.google.gwt.user.client.rpc.IsSerializable;
public class ReplicaDTO implements IsSerializable {
private String hostname;
private String identifier;
private Date registrationTime;
private Map<String, Integer> operationCountByOperationClassName;
ReplicaDTO() {}
public ReplicaDTO(String hostname, Date registrationTime,
Map<String, Integer> operationCountByOperationClassName) {
String identifier, Map<String, Integer> operationCountByOperationClassName) {
this.hostname = hostname;
this.identifier = identifier;
this.registrationTime = registrationTime;
this.operationCountByOperationClassName = operationCountByOperationClassName;
}
@@ -25,4 +27,7 @@ public class ReplicaDTO implements IsSerializable {
public Map<String, Integer> getOperationCountByOperationClassName() {
return operationCountByOperationClassName;
}
public String getIdentifier() {
return identifier;
}
}
@@ -14,14 +14,16 @@ import com.google.gwt.user.client.rpc.IsSerializable;
public class ReplicationStateDTO implements IsSerializable {
private Map<String, ReplicaDTO> replicaInfoByHostname;
private ReplicationMasterDTO replicatingFromMaster;
private String serverIdentifier;
ReplicationStateDTO() { } // for de-serialization
public ReplicationStateDTO(ReplicationMasterDTO replicatingFromMaster, Iterable<ReplicaDTO> replicas) {
public ReplicationStateDTO(ReplicationMasterDTO replicatingFromMaster, Iterable<ReplicaDTO> replicas, String serverIdentifier) {
this.replicatingFromMaster = replicatingFromMaster;
this.serverIdentifier = serverIdentifier;
this.replicaInfoByHostname = new HashMap<String, ReplicaDTO>();
for (ReplicaDTO replica : replicas) {
replicaInfoByHostname.put(replica.getHostname(), replica);
replicaInfoByHostname.put(replica.getIdentifier(), replica);
}
}
@@ -33,7 +35,7 @@ public class ReplicationStateDTO implements IsSerializable {
return replicaInfoByHostname.values();
}
public ReplicaDTO getReplicaByHostname(String hostname) {
return replicaInfoByHostname.get(hostname);
public String getServerIdentifier() {
return serverIdentifier;
}
}
@@ -25,9 +25,9 @@ import com.sap.sailing.domain.common.impl.Util.Pair;
import com.sap.sailing.mongodb.MongoDBService;
import com.sap.sailing.server.RacingEventService;
import com.sap.sailing.server.impl.RacingEventServiceImpl;
import com.sap.sailing.server.replication.ReplicaDescriptor;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
import com.sap.sailing.server.replication.ReplicationService;
import com.sap.sailing.server.replication.impl.ReplicaDescriptor;
import com.sap.sailing.server.replication.impl.ReplicationInstancesManager;
import com.sap.sailing.server.replication.impl.ReplicationMasterDescriptorImpl;
import com.sap.sailing.server.replication.impl.ReplicationServiceImpl;
@@ -38,6 +38,7 @@ public abstract class AbstractServerReplicationTest {
private DomainFactory resolveAgainst;
protected RacingEventServiceImpl replica;
protected RacingEventServiceImpl master;
protected ReplicationServiceTestImpl replicaReplicator;
private ReplicaDescriptor replicaDescriptor;
private ReplicationServiceImpl masterReplicator;
@@ -70,6 +71,7 @@ public abstract class AbstractServerReplicationTest {
protected Pair<ReplicationServiceTestImpl, ReplicationMasterDescriptor> basicSetUp(
boolean dropDB, RacingEventServiceImpl master, RacingEventServiceImpl replica) throws IOException, InterruptedException {
final String exchangeName = "test-sapsailinganalytics-exchange";
final UUID serverUuid = UUID.randomUUID();
final MongoDBService mongoDBService = MongoDBService.INSTANCE;
if (dropDB) {
mongoDBService.getDB().dropDatabase();
@@ -87,13 +89,14 @@ public abstract class AbstractServerReplicationTest {
}
ReplicationInstancesManager rim = new ReplicationInstancesManager();
masterReplicator = new ReplicationServiceImpl(exchangeName, rim, this.master);
replicaDescriptor = new ReplicaDescriptor(InetAddress.getLocalHost());
replicaDescriptor = new ReplicaDescriptor(InetAddress.getLocalHost(), serverUuid, "");
masterReplicator.registerReplica(replicaDescriptor);
ReplicationMasterDescriptor masterDescriptor = new ReplicationMasterDescriptorImpl("localhost", exchangeName, SERVLET_PORT, 0);
ReplicationMasterDescriptor masterDescriptor = new ReplicationMasterDescriptorImpl("localhost", exchangeName, SERVLET_PORT, 0, "test-queue");
ReplicationServiceTestImpl replicaReplicator = new ReplicationServiceTestImpl(exchangeName, resolveAgainst, rim,
replicaDescriptor, this.replica, this.master, masterReplicator, masterDescriptor);
Pair<ReplicationServiceTestImpl, ReplicationMasterDescriptor> result = new Pair<>(replicaReplicator, masterDescriptor);
replicaReplicator.startInitialLoadTransmissionServlet();
this.replicaReplicator = replicaReplicator;
return result;
}
@@ -103,7 +106,11 @@ public abstract class AbstractServerReplicationTest {
URLConnection urlConnection = new URL("http://localhost:"+SERVLET_PORT+"/STOP").openConnection(); // stop the initial load test server thread
urlConnection.getInputStream().close();
}
public void stopReplicatingToMaster() throws IOException {
replicaReplicator.stopToReplicateFromMaster();
}
static class ReplicationServiceTestImpl extends ReplicationServiceImpl {
private final DomainFactory resolveAgainst;
private final RacingEventService master;
@@ -144,7 +151,12 @@ public abstract class AbstractServerReplicationTest {
pw.println("Content-Type: text/plain");
pw.println();
pw.flush();
if (request.contains("REGISTER")) {
if (request.contains("DEREGISTER")) {
// assuming that it is safe to unregister all replicas for tests
for (ReplicaDescriptor descriptor : getReplicaInfo()) {
unregisterReplica(descriptor);
}
} else if (request.contains("REGISTER")) {
final String uuid = UUID.randomUUID().toString();
registerReplicaUuidForMaster(uuid, masterDescriptor);
pw.print(uuid.getBytes());
@@ -0,0 +1,132 @@
package com.sap.sailing.server.replication.test;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNotSame;
import static org.junit.Assert.assertNull;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import java.util.logging.Logger;
import org.junit.Before;
import org.junit.Test;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.ConsumerCancelledException;
import com.rabbitmq.client.QueueingConsumer;
import com.rabbitmq.client.ShutdownSignalException;
import com.sap.sailing.domain.base.Event;
import com.sap.sailing.domain.common.impl.Util;
import com.sap.sailing.domain.common.impl.Util.Pair;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
import com.sap.sailing.server.replication.impl.ReplicationMasterDescriptorImpl;
public class ConnectionResetAndReconnectTest extends AbstractServerReplicationTest {
static final Logger logger = Logger.getLogger(ConnectionResetAndReconnectTest.class.getName());
public static boolean forceStopDelivery = false;
static class QueuingConsumerTest extends QueueingConsumer {
public QueuingConsumerTest(Channel ch) {
super(ch);
}
@Override
public Delivery nextDelivery() throws ShutdownSignalException, ConsumerCancelledException, InterruptedException {
if (forceStopDelivery) {
throw new ShutdownSignalException(false, false, null, null);
}
return super.nextDelivery();
}
}
static class MasterReplicationDescriptorMock extends ReplicationMasterDescriptorImpl {
public MasterReplicationDescriptorMock(String hostname, String exchangeName, int servletPort, int messagingPort) {
super(hostname, exchangeName, servletPort, messagingPort, "dummy");
}
public static MasterReplicationDescriptorMock from(ReplicationMasterDescriptor obj) {
return new MasterReplicationDescriptorMock(obj.getHostname(), obj.getExchangeName(), obj.getServletPort(), obj.getMessagingPort());
}
@Override
public QueueingConsumer getConsumer() throws IOException {
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost(getHostname());
int port = getMessagingPort();
if (port != 0) {
connectionFactory.setPort(port);
}
Connection connection = connectionFactory.newConnection();
Channel channel = connection.createChannel();
channel.exchangeDeclare(getExchangeName(), "fanout");
QueueingConsumer consumer = new QueuingConsumerTest(channel);
String queueName = channel.queueDeclare().getQueue();
channel.queueBind(queueName, getExchangeName(), "");
channel.basicConsume(queueName, /* auto-ack */ true, consumer);
return consumer;
}
}
private MasterReplicationDescriptorMock masterReplicationDescriptor;
private ReplicationServiceTestImpl replicaReplicationDescriptor;
@Before
public void setUp() throws Exception {
try {
Pair<ReplicationServiceTestImpl, ReplicationMasterDescriptor> result = basicSetUp(
true, /* master=null means create a new one */ null,
/* replica=null means create a new one */null);
masterReplicationDescriptor = MasterReplicationDescriptorMock.from(result.getB());
replicaReplicationDescriptor = result.getA();
} catch (Exception e) {
e.printStackTrace();
tearDown();
}
}
@Test
public void testReplicaLoosingConnectionToExchangeQueue() throws Exception {
assertNotSame(master, replica);
assertEquals(Util.size(master.getAllRegattas()), Util.size(replica.getAllRegattas()));
/* until here both instances should have the same in-memory state.
* now lets add an event on master and stop the messaging queue. */
stopMessagingExchange();
replicaReplicationDescriptor.startToReplicateFrom(masterReplicationDescriptor);
Event event = addEventOnMaster();
Thread.sleep(1000); // wait for master queue to get filled
assertNull(replica.getEvent(event.getId()));
startMessagingExchange();
Thread.sleep(3000); // wait for connection to recover
assertNotNull(replica.getEvent(event.getId()));
}
private Event addEventOnMaster() {
final String eventName = "ESS Masquat";
final String venueName = "Masquat, Oman";
final String publicationUrl = "http://ess40.sapsailing.com";
final boolean isPublic = false;
List<String> regattas = new ArrayList<String>();
regattas.add("Day1");
regattas.add("Day2");
return master.addEvent(eventName, venueName, publicationUrl, isPublic, UUID.randomUUID(), regattas);
}
private void stopMessagingExchange() {
forceStopDelivery = true;
}
private void startMessagingExchange() {
forceStopDelivery = false;
}
}
@@ -0,0 +1,40 @@
package com.sap.sailing.server.replication.test;
import static org.junit.Assert.assertEquals;
import java.net.InetAddress;
import java.net.UnknownHostException;
import java.util.Arrays;
import java.util.UUID;
import org.junit.Before;
import org.junit.Test;
import com.sap.sailing.server.operationaltransformation.CreateLeaderboardGroup;
import com.sap.sailing.server.replication.impl.ReplicaDescriptor;
import com.sap.sailing.server.replication.impl.ReplicationInstancesManager;
public class ReplicationInstancesManagerLoggingPerformanceTest {
private ReplicationInstancesManager replicationInstanceManager;
private CreateLeaderboardGroup operation;
private ReplicaDescriptor replica;
@Before
public void setUp() throws UnknownHostException {
replicationInstanceManager = new ReplicationInstancesManager();
replica = new ReplicaDescriptor(InetAddress.getLocalHost(), UUID.randomUUID(), "");
replicationInstanceManager.registerReplica(replica);
operation = new CreateLeaderboardGroup("Test Leaderboard Group", "Description of Test Leaderboard Group", /* displayGroupsInReverseOrder */ false, Arrays.asList(new String[] { "Default Leaderboard" }),
/* overallLeaderboardDiscardThresholds */ null, /* overallLeaderboardScoringSchemeType */ null);
}
@Test
public void testLoggingPerformance() {
final int count = 10000000;
for (int i=0; i<count; i++) {
replicationInstanceManager.log(operation);
}
assertEquals(count, replicationInstanceManager.getStatistics(replica).get(operation.getClass()).intValue());
}
}
@@ -0,0 +1,41 @@
package com.sap.sailing.server.replication.test;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNotSame;
import static org.junit.Assert.assertNull;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import org.junit.Test;
import com.sap.sailing.domain.base.Event;
import com.sap.sailing.domain.common.impl.Util;
public class ReplicationStopTest extends AbstractServerReplicationTest {
@Test
public void testReplicaStop() throws Exception {
assertNotSame(master, replica);
assertEquals(Util.size(master.getAllRegattas()), Util.size(replica.getAllRegattas()));
Event event = addEventOnMaster("Test 1 Event");
Thread.sleep(1000);
assertNotNull(replica.getEvent(event.getId()));
stopReplicatingToMaster();
Event event2 = addEventOnMaster("Test 2 Event");
assertNull(replica.getEvent(event2.getId()));
}
private Event addEventOnMaster(String eventName) {
final String venueName = "Masquat, Oman";
final String publicationUrl = "http://ess40.sapsailing.com";
final boolean isPublic = false;
List<String> regattas = new ArrayList<String>();
regattas.add("Day1");
regattas.add("Day2");
return master.addEvent(eventName, venueName, publicationUrl, isPublic, UUID.randomUUID(), regattas);
}
}
@@ -7,7 +7,8 @@ Bundle-Activator: com.sap.sailing.server.replication.impl.Activator
Bundle-Vendor: SAP
Bundle-RequiredExecutionEnvironment: JavaSE-1.7
Bundle-ActivationPolicy: lazy
Export-Package: com.sap.sailing.server.replication
Export-Package: com.sap.sailing.server.replication,
com.sap.sailing.server.replication.impl
Require-Bundle: com.sap.sailing.server,
com.sap.sailing.server.gateway,
com.sap.sailing.domain.common,
@@ -5,5 +5,5 @@ import com.sap.sailing.server.replication.impl.ReplicationFactoryImpl;
public interface ReplicationFactory {
static ReplicationFactory INSTANCE = new ReplicationFactoryImpl();
ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int jmsPort);
ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int jmsPort, String jmsQueueName);
}
@@ -1,8 +1,10 @@
package com.sap.sailing.server.replication;
import java.io.IOException;
import java.io.UnsupportedEncodingException;
import java.net.MalformedURLException;
import java.net.URL;
import java.util.UUID;
import com.rabbitmq.client.QueueingConsumer;
@@ -14,7 +16,9 @@ import com.rabbitmq.client.QueueingConsumer;
*/
public interface ReplicationMasterDescriptor {
URL getReplicationRegistrationRequestURL() throws MalformedURLException;
URL getReplicationRegistrationRequestURL(UUID uuid, String additionalInformation) throws MalformedURLException, UnsupportedEncodingException;
URL getReplicationDeRegistrationRequestURL(UUID uuid) throws MalformedURLException;
URL getInitialLoadURL() throws MalformedURLException;
@@ -32,4 +36,6 @@ public interface ReplicationMasterDescriptor {
QueueingConsumer getConsumer() throws IOException;
String getExchangeName();
void stopConnection();
}
@@ -2,8 +2,10 @@ package com.sap.sailing.server.replication;
import java.io.IOException;
import java.util.Map;
import java.util.UUID;
import com.sap.sailing.server.RacingEventServiceOperation;
import com.sap.sailing.server.replication.impl.ReplicaDescriptor;
import com.sap.sailing.server.replication.impl.ReplicationServlet;
public interface ReplicationService {
@@ -41,4 +43,23 @@ public interface ReplicationService {
* replica by type, where the operation type is the key, represented as the operation's class name
*/
Map<Class<? extends RacingEventServiceOperation<?>>, Integer> getStatistics(ReplicaDescriptor replicaDescriptor);
/**
* Stops the currently running replication. As there can be only one replication running
* this method needs no parameters.
* @throws IOException
*/
void stopToReplicateFromMaster() throws IOException;
/**
* Stops all replica currently registered with this server.
* @throws IOException
*/
void stopAllReplica() throws IOException;
/**
* Returns an unique server identifier
* @return
*/
UUID getServerIdentifier();
}
@@ -1,4 +1,4 @@
package com.sap.sailing.server.replication;
package com.sap.sailing.server.replication.impl;
import java.io.Serializable;
import java.net.InetAddress;
@@ -18,18 +18,18 @@ public class ReplicaDescriptor implements Serializable {
private static final long serialVersionUID = -5451556877949921454L;
private final UUID uuid;
private final InetAddress ipAddress;
private final TimePoint registrationTime;
private final String additionalInformation;
/**
* Sets the registration time to now.
*/
public ReplicaDescriptor(InetAddress ipAddress) {
this.uuid = UUID.randomUUID();
public ReplicaDescriptor(InetAddress ipAddress, UUID serverUuid, String additionalInformation) {
this.uuid = serverUuid;
this.registrationTime = MillisecondsTimePoint.now();
this.ipAddress = ipAddress;
this.additionalInformation = additionalInformation;
}
public UUID getUuid() {
@@ -43,6 +43,10 @@ public class ReplicaDescriptor implements Serializable {
public TimePoint getRegistrationTime() {
return registrationTime;
}
public String getAdditionalInformation() {
return additionalInformation;
}
@Override
public int hashCode() {
@@ -5,8 +5,8 @@ import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
public class ReplicationFactoryImpl implements ReplicationFactory {
@Override
public ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int messagingPort) {
return new ReplicationMasterDescriptorImpl(hostname, exchangeName, servletPort, messagingPort);
public ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int messagingPort, String queueName) {
return new ReplicationMasterDescriptorImpl(hostname, exchangeName, servletPort, messagingPort, queueName);
}
}
@@ -7,7 +7,6 @@ import java.util.Map;
import java.util.Set;
import com.sap.sailing.server.RacingEventServiceOperation;
import com.sap.sailing.server.replication.ReplicaDescriptor;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
public class ReplicationInstancesManager {
@@ -86,4 +85,9 @@ public class ReplicationInstancesManager {
public Map<Class<? extends RacingEventServiceOperation<?>>, Integer> getStatistics(ReplicaDescriptor replica) {
return replicationCounts.get(replica);
}
public void removeAll() {
replicationCounts.clear();
replicaDescriptors.clear();
}
}
@@ -1,14 +1,19 @@
package com.sap.sailing.server.replication.impl;
import java.io.IOException;
import java.io.UnsupportedEncodingException;
import java.net.MalformedURLException;
import java.net.URL;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.QueueingConsumer;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
import com.sap.sailing.util.BuildVersion;
public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescriptor {
private static final String REPLICATION_SERVLET = "/replication/replication";
@@ -16,23 +21,37 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
private final String exchangeName;
private final int servletPort;
private final int messagingPort;
private final String queueName;
private QueueingConsumer consumer;
/**
* @param messagingPort 0 means use default port
*/
public ReplicationMasterDescriptorImpl(String hostname, String exchangeName, int servletPort, int messagingPort) {
public ReplicationMasterDescriptorImpl(String hostname, String exchangeName, int servletPort, int messagingPort, String queueName) {
this.hostname = hostname;
this.servletPort = servletPort;
this.messagingPort = messagingPort;
this.exchangeName = exchangeName;
this.queueName = queueName;
this.consumer = null;
}
@Override
public URL getReplicationRegistrationRequestURL() throws MalformedURLException {
public URL getReplicationRegistrationRequestURL(UUID uuid, String additional) throws MalformedURLException, UnsupportedEncodingException {
return new URL("http", hostname, servletPort, REPLICATION_SERVLET + "?" + ReplicationServlet.ACTION + "="
+ ReplicationServlet.Action.REGISTER.name());
+ ReplicationServlet.Action.REGISTER.name()
+ "&" + ReplicationServlet.SERVER_UUID + "=" + uuid.toString()
+ "&" + ReplicationServlet.ADDITIONAL_INFORMATION + "=" + java.net.URLEncoder.encode(BuildVersion.getBuildVersion(), "UTF-8"));
}
@Override
public URL getReplicationDeRegistrationRequestURL(UUID uuid) throws MalformedURLException {
return new URL("http", hostname, servletPort, REPLICATION_SERVLET + "?" + ReplicationServlet.ACTION + "="
+ ReplicationServlet.Action.DEREGISTER.name()
+ "&" + ReplicationServlet.SERVER_UUID + "=" + uuid.toString());
}
@Override
public URL getInitialLoadURL() throws MalformedURLException {
return new URL("http", hostname, servletPort, REPLICATION_SERVLET + "?" + ReplicationServlet.ACTION + "="
@@ -40,7 +59,7 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
}
@Override
public QueueingConsumer getConsumer() throws IOException {
public synchronized QueueingConsumer getConsumer() throws IOException {
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost(getHostname());
int port = getMessagingPort();
@@ -49,13 +68,67 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
}
Connection connection = connectionFactory.newConnection();
Channel channel = connection.createChannel();
/*
* Connect a queue to the given exchange that has already
* been created by the master server.
*/
channel.exchangeDeclare(exchangeName, "fanout");
QueueingConsumer consumer = new QueueingConsumer(channel);
String queueName = channel.queueDeclare().getQueue();
/*
* The x-message-ttl argument to queue.declare controls for how long a message published to a queue can live before
* it is discarded. A message that has been in the queue for longer than the configured TTL is said to be dead.
* Note that a message routed to multiple queues can die at different times, or not at all,
* in each queue in which it resides. The death of a message in one queue has no impact on the life of the
* same message in other queues.
*/
final Map<String, Object> args = new HashMap<String, Object>();
args.put("x-message-ttl", (60*30)*1000); // messages will live half an hour in queue before being deleted
/*
* The x-expires argument to queue.declare controls for how long a queue can be unused before it is automatically
* deleted. Unused means the queue has no consumers, the queue has not been redeclared, and basic.get has not
* been invoked for a duration of at least the expiration period.
*/
args.put("x-expires", (60*60)*1000); // queue will live one hour before being deleted
/*
* The maximum length of a queue can be limited to a set number of messages by supplying the x-max-length queue
* declaration argument with a non-negative integer value. Queue length is a measure that takes into account
* ready messages, ignoring unacknowledged messages and message size. Messages will be dropped or dead-lettered
* from the front of the queue to make room for new messages once the limit is reached.
*/
args.put("x-max-length", 3000000);
// a server-named non-exclusive, non-durable queue
// this queue will survive a connection drop (autodelete=false) and
// will also support being reconnected (exclusive=false). it will
// not survive a rabbitmq server restart (durable=false).
String queueName = channel.queueDeclare(this.queueName,
/*durable*/ false, /*exclusive*/ false, /*auto-delete*/ false, args).getQueue();
// from now on we get all new messages that the exchange is getting from producer
channel.queueBind(queueName, exchangeName, "");
channel.basicConsume(queueName, /* auto-ack */ true, consumer);
this.consumer = consumer;
return consumer;
}
@Override
public synchronized void stopConnection() {
try {
if (consumer != null) {
// make sure to remove queue in order to avoid any exchanges filling it with messages
consumer.getChannel().queueUnbind(queueName, exchangeName, "");
consumer.getChannel().queueDelete(queueName);
consumer.getChannel().getConnection().close(1);
}
} catch (Exception ex) {
// ignore any exception during abort. close can yield a broad
// number of exceptions that we don't want to know or to log.
}
}
/**
* @return 0 means use default port
@@ -79,4 +152,9 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
public String getExchangeName() {
return exchangeName;
}
public String toString() {
return getHostname() + ":" + getServletPort() + "/" + getMessagingPort();
}
}
@@ -5,10 +5,13 @@ import java.io.IOException;
import java.io.InputStream;
import java.io.ObjectInputStream;
import java.io.ObjectOutputStream;
import java.net.ConnectException;
import java.net.URL;
import java.net.URLConnection;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import java.util.logging.Level;
import java.util.logging.Logger;
import org.osgi.util.tracker.ServiceTracker;
@@ -20,9 +23,9 @@ import com.rabbitmq.client.QueueingConsumer;
import com.sap.sailing.server.OperationExecutionListener;
import com.sap.sailing.server.RacingEventService;
import com.sap.sailing.server.RacingEventServiceOperation;
import com.sap.sailing.server.replication.ReplicaDescriptor;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
import com.sap.sailing.server.replication.ReplicationService;
import com.sap.sailing.util.BuildVersion;
/**
* Can observe a {@link RacingEventService} for the operations it performs that require replication. Only observes as
@@ -69,6 +72,14 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
*/
private final String exchangeName;
/**
* UUID that identifies this server
*/
private final UUID serverUUID;
private Replicator replicator;
private Thread replicatorThread;
public ReplicationServiceImpl(String exchangeName, final ReplicationInstancesManager replicationInstancesManager) throws IOException {
this.replicationInstancesManager = replicationInstancesManager;
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>();
@@ -77,6 +88,9 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
racingEventServiceTracker.open();
localService = null;
this.exchangeName = exchangeName;
replicator = null;
serverUUID = UUID.randomUUID();
logger.info("Setting " + serverUUID.toString() + " as unique replication identifier.");
}
/**
@@ -87,15 +101,28 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
public ReplicationServiceImpl(String exchangeName,
final ReplicationInstancesManager replicationInstancesManager, RacingEventService localService) throws IOException {
this.replicationInstancesManager = replicationInstancesManager;
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>();
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;
replicator = null;
serverUUID = UUID.randomUUID();
logger.info("Setting " + serverUUID.toString() + " as unique replication identifier.");
}
private Channel createChannel(String exchangeName) throws IOException {
private Channel createMasterChannel(String exchangeName) throws IOException {
final ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("localhost"); // ...and use default port
Channel result = connectionFactory.newConnection().createChannel();
Channel result = null;
try {
result = connectionFactory.newConnection().createChannel();
} catch (ConnectException ex) {
// make sure to log something meaningful
logger.severe("Could not connect to messaging queue on " + connectionFactory.getHost() + ":" + connectionFactory.getPort() + "/" + exchangeName);
throw ex;
}
logger.info("Connected to " + connectionFactory.getHost() + ":" + connectionFactory.getPort() + "/" + exchangeName);
result.exchangeDeclare(exchangeName, "fanout");
return result;
}
@@ -117,7 +144,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
addAsListenerToRacingEventService();
synchronized (this) {
if (masterChannel == null) {
masterChannel = createChannel(exchangeName);
masterChannel = createMasterChannel(exchangeName);
}
}
}
@@ -134,8 +161,10 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
if (!replicationInstancesManager.hasReplicas()) {
removeAsListenerFromRacingEventService();
synchronized (this) {
masterChannel.close();
masterChannel = null;
if (masterChannel != null) {
masterChannel.close();
masterChannel = null;
}
}
}
}
@@ -169,14 +198,31 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
@Override
public void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException, InterruptedException {
logger.info("Starting to replicate from "+master);
try {
registerReplicaWithMaster(master);
} catch (Exception ex) {
logger.log(Level.SEVERE, "ERROR", ex);
throw ex;
}
replicatingFromMaster = master;
registerReplicaWithMaster(master);
QueueingConsumer consumer = master.getConsumer();
logger.info("Registered replica with master");
QueueingConsumer consumer = null;
// logging exception here because it will not propagate
// thru the client with all details
try {
consumer = master.getConsumer();
} catch (Exception ex) {
logger.log(Level.SEVERE, "ERROR", ex);
replicatingFromMaster = null;
throw ex;
}
logger.info("Connection to exchange successful.");
URL initialLoadURL = master.getInitialLoadURL();
logger.info("Initial load URL is "+initialLoadURL);
final Replicator replicator = new Replicator(master, this, /* startSuspended */ true, consumer);
replicator = new Replicator(master, this, /* startSuspended */ true, consumer);
// start receiving messages already now, but start in suspended mode
new Thread(replicator, "Replicator receiving from "+master.getHostname()+"/"+master.getExchangeName()).start();
replicatorThread = new Thread(replicator, "Replicator receiving from "+master.getHostname()+"/"+master.getExchangeName());
replicatorThread.start();
logger.info("Started replicator thread");
InputStream is = initialLoadURL.openStream();
final RacingEventService racingEventService = getRacingEventService();
@@ -191,7 +237,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
* @return the UUID that the master generated for this client which is also entered into {@link #replicaUUIDs}
*/
private String registerReplicaWithMaster(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException {
URL replicationRegistrationRequestURL = master.getReplicationRegistrationRequestURL();
URL replicationRegistrationRequestURL = master.getReplicationRegistrationRequestURL(getServerIdentifier(), BuildVersion.getBuildVersion());
final URLConnection registrationRequestConnection = replicationRegistrationRequestURL.openConnection();
registrationRequestConnection.connect();
InputStream content = (InputStream) registrationRequestConnection.getContent();
@@ -206,6 +252,27 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
registerReplicaUuidForMaster(replicaUUID, master);
return replicaUUID;
}
protected void deregisterReplicaWithMaster(ReplicationMasterDescriptor master) {
try {
URL replicationDeRegistrationRequestURL = master.getReplicationDeRegistrationRequestURL(getServerIdentifier());
final URLConnection deregistrationRequestConnection = replicationDeRegistrationRequestURL.openConnection();
deregistrationRequestConnection.connect();
StringBuilder uuid = new StringBuilder();
InputStream content = (InputStream) deregistrationRequestConnection.getContent();
byte[] buf = new byte[256];
int read = content.read(buf);
while (read != -1) {
uuid.append(new String(buf, 0, read));
read = content.read(buf);
}
content.close();
} catch (Exception ex) {
// ignore exceptions here - they will mostly be caused by an incompatible server
// it is also not problematic if the server does not get this deregistration
// a new registration will overwrite the current one
}
}
@Override
public <T> void executed(RacingEventServiceOperation<T> operation) {
@@ -225,4 +292,44 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
return replicationInstancesManager.getStatistics(replicaDescriptor);
}
@Override
public void stopToReplicateFromMaster() throws IOException {
ReplicationMasterDescriptor descriptor = isReplicatingFromMaster();
if (descriptor != null) {
synchronized(replicaUUIDs) {
if (replicator != null) {
replicator.stop();
deregisterReplicaWithMaster(descriptor);
replicatingFromMaster = null;
replicaUUIDs.clear();
// this is needed because QueuingConsumer.nextDelivery() wont unblock
// if the connection is closed by application.
replicatorThread.interrupt();
replicator = null;
}
}
}
}
@Override
public void stopAllReplica() throws IOException {
if (replicationInstancesManager.hasReplicas()) {
replicationInstancesManager.removeAll();
removeAsListenerFromRacingEventService();
synchronized (this) {
if (masterChannel != null) {
masterChannel.close();
masterChannel = null;
}
}
logger.info("Unregistered all replicas from this server!");
}
}
@Override
public UUID getServerIdentifier() {
return serverUUID;
}
}
@@ -5,6 +5,7 @@ import java.io.ObjectOutputStream;
import java.net.InetAddress;
import java.net.UnknownHostException;
import java.util.Arrays;
import java.util.UUID;
import java.util.logging.Level;
import java.util.logging.Logger;
@@ -17,7 +18,6 @@ import org.osgi.util.tracker.ServiceTracker;
import com.sap.sailing.server.RacingEventService;
import com.sap.sailing.server.gateway.SailingServerHttpServlet;
import com.sap.sailing.server.replication.ReplicaDescriptor;
import com.sap.sailing.server.replication.ReplicationService;
/**
@@ -32,9 +32,11 @@ public class ReplicationServlet extends SailingServerHttpServlet {
private static final long serialVersionUID = 4835516998934433846L;
public enum Action { REGISTER, INITIAL_LOAD }
public enum Action { REGISTER, INITIAL_LOAD, DEREGISTER }
public static final String ACTION = "action";
public static final String SERVER_UUID = "uuid";
public static final String ADDITIONAL_INFORMATION = "additional";
private ServiceTracker<ReplicationService, ReplicationService> replicationServiceTracker;
@@ -62,6 +64,9 @@ public class ReplicationServlet extends SailingServerHttpServlet {
case REGISTER:
registerClientWithReplicationService(req, resp);
break;
case DEREGISTER:
deregisterClientWithReplicationService(req, resp);
break;
case INITIAL_LOAD:
ObjectOutputStream oos = new ObjectOutputStream(resp.getOutputStream());
try {
@@ -79,6 +84,14 @@ public class ReplicationServlet extends SailingServerHttpServlet {
}
}
private void deregisterClientWithReplicationService(HttpServletRequest req, HttpServletResponse resp) throws IOException {
ReplicaDescriptor replica = getReplicaDescriptor(req);
getReplicationService().unregisterReplica(replica);
logger.info("Deregistered replication client with this server " + replica.getIpAddress());
resp.setContentType("text/plain");
resp.getWriter().print(replica.getUuid());
}
private void registerClientWithReplicationService(HttpServletRequest req, HttpServletResponse resp)
throws IOException {
ReplicaDescriptor replica = getReplicaDescriptor(req);
@@ -89,6 +102,9 @@ public class ReplicationServlet extends SailingServerHttpServlet {
private ReplicaDescriptor getReplicaDescriptor(HttpServletRequest req) throws UnknownHostException {
InetAddress ipAddress = InetAddress.getByName(req.getRemoteAddr());
return new ReplicaDescriptor(ipAddress);
UUID uuid = UUID.fromString(req.getParameter(SERVER_UUID));
String additional = req.getParameter(ADDITIONAL_INFORMATION);
logger.info("Registered new replica " + ipAddress + " " + uuid.toString() + " " + additional);
return new ReplicaDescriptor(ipAddress, uuid, additional);
}
}
@@ -1,6 +1,7 @@
package com.sap.sailing.server.replication.impl;
import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.io.ObjectInputStream;
import java.util.ArrayList;
import java.util.Iterator;
@@ -32,16 +33,27 @@ import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
public class Replicator implements Runnable {
private final static Logger logger = Logger.getLogger(Replicator.class.getName());
private static final long CHECK_INTERVAL_MILLIS = 2000; // how long (milliseconds) to pause before checking connection again
private static final int CHECK_COUNT = 150; // how long to check, value is CHECK_INTERVAL second steps
private final ReplicationMasterDescriptor master;
private final HasRacingEventService racingEventServiceTracker;
private final List<RacingEventServiceOperation<?>> queue;
private final QueueingConsumer consumer;
private QueueingConsumer consumer;
/**
* How many checks have been performed due to a failing connection?
*/
private int checksPerformed = 0;
/**
* If the replicator is suspended, messages received are queued.
*/
private boolean suspended;
private boolean stopped = false;
/**
* Starts the replicator immediately, not holding back messages received but forwarding them directly.
*
@@ -71,10 +83,15 @@ public class Replicator implements Runnable {
@Override
public void run() {
ClassLoader oldClassLoader = Thread.currentThread().getContextClassLoader();
while (true) {
if (isBeingStopped()) {
break;
}
try {
Delivery delivery = consumer.nextDelivery();
byte[] bytesFromMessage = delivery.getBody();
checksPerformed = 0;
// Set this object's class's class loader as context for de-serialization so that all exported classes
// of all required bundles/packages can be deserialized at least
Thread.currentThread().setContextClassLoader(getClass().getClassLoader());
@@ -82,9 +99,49 @@ public class Replicator implements Runnable {
new ByteArrayInputStream(bytesFromMessage));
RacingEventServiceOperation<?> operation = (RacingEventServiceOperation<?>) ois.readObject();
applyOrQueue(operation);
} catch (InterruptedException irr) {
logger.info("Application requested shutdown.");
} catch (ShutdownSignalException sse) {
logger.info("Received "+sse.getMessage()+". Terminating "+this);
break;
/* make sure to respond to a stop event without waiting */
if (isBeingStopped()) {
break;
}
if (sse.isInitiatedByApplication()) {
logger.severe("Application shut down messaging queue for " + this.toString());
break;
}
logger.info(sse.getMessage());
if (checksPerformed <= CHECK_COUNT) {
try {
Thread.sleep(CHECK_INTERVAL_MILLIS);
/* isOpen() will return false if the channel has been closed. This
* does not hold when the connection is dropped.
*/
if (!this.consumer.getChannel().isOpen()) {
/* for a reconnection we need to instantiate a new consumer */
try {
logger.info("Channel seems to be closed. Trying to reconnect consumer queue...");
this.consumer = master.getConsumer();
logger.info("OK - channel reconnected!");
Thread.sleep(CHECK_INTERVAL_MILLIS);
checksPerformed += 1;
} catch (IOException eio) {
// do not print exceptions known to occur
}
}
} catch (InterruptedException eir) {
eir.printStackTrace();
}
checksPerformed += 1;
continue;
} else {
logger.severe("Grace time (" + CHECK_COUNT*(CHECK_INTERVAL_MILLIS/1000) + "secs) is over. Terminating replication listener " + this.toString());
// XXX: Also make sure that all handlers get notifications about this
break;
}
} catch (Exception e) {
logger.info("Exception while processing replica: "+e.getMessage());
logger.log(Level.SEVERE, "run", e);
@@ -92,6 +149,8 @@ public class Replicator implements Runnable {
Thread.currentThread().setContextClassLoader(oldClassLoader);
}
}
logger.info("Stopped replicator thread. This server will no longer receive events from a master.");
}
public synchronized boolean isQueueEmpty() {
@@ -144,6 +203,20 @@ public class Replicator implements Runnable {
public synchronized boolean isSuspended() {
return suspended;
}
public synchronized void stop() {
if (isSuspended()) {
/* make sure to apply everything in queue before stopping this thread */
applyQueue();
}
logger.info("Signaled Replicator thread to stop asap.");
stopped = true;
master.stopConnection();
}
public synchronized boolean isBeingStopped() {
return stopped;
}
@Override
public String toString() {
@@ -145,6 +145,7 @@ import com.sap.sailing.server.operationaltransformation.UpdateSpecificRegatta;
import com.sap.sailing.server.operationaltransformation.UpdateTrackedRaceStatus;
import com.sap.sailing.server.operationaltransformation.UpdateWindAveragingTime;
import com.sap.sailing.server.operationaltransformation.UpdateWindSourcesToExclude;
import com.sap.sailing.util.BuildVersion;
public class RacingEventServiceImpl implements RacingEventService, RegattaListener, LeaderboardRegistry, Replicator {
private static final Logger logger = Logger.getLogger(RacingEventServiceImpl.class.getName());
@@ -1645,26 +1646,55 @@ public class RacingEventServiceImpl implements RacingEventService, RegattaListen
@Override
public void serializeForInitialReplication(ObjectOutputStream oos) throws IOException {
logger.info("serializing eventsById");
StringBuffer logoutput = new StringBuffer();
logger.info("Serializing events...");
oos.writeObject(eventsById);
logger.info("serializing regattasByName");
logoutput.append("\nSerialized " + eventsById.size() + " events\n");
for (Event event : eventsById.values()) {
logoutput.append(String.format("%3s\n", event.toString()));
}
logger.info("Serializing regattas...");
oos.writeObject(regattasByName);
logger.info("serializing regattasObservedForDefaultLeaderboard");
logoutput.append("Serialized " + regattasByName.size() + " regattas\n");
for (Regatta regatta : regattasByName.values()) {
logoutput.append(String.format("%3s\n", regatta.toString()));
}
logger.info("Serializing regattas observed...");
oos.writeObject(regattasObservedForDefaultLeaderboard);
logger.info("serializing regattaTrackingCache");
logger.info("Serializing regatta tracking cache...");
oos.writeObject(regattaTrackingCache);
logger.info("serializing leaderboardGroupsByName");
logger.info("Serializing leaderboard groups...");
oos.writeObject(leaderboardGroupsByName);
logger.info("serializing leaderboardsByName");
logoutput.append("Serialized " + leaderboardGroupsByName.size() + " leaderboard groups\n");
for (LeaderboardGroup lg : leaderboardGroupsByName.values()) {
logoutput.append(String.format("%3s\n", lg.toString()));
}
logger.info("Serializing leaderboards...");
oos.writeObject(leaderboardsByName);
logger.info("serializing mediaLibrary");
logoutput.append("Serialized " + leaderboardsByName.size() + " leaderboards\n");
for (Leaderboard lg : leaderboardsByName.values()) {
logoutput.append(String.format("%3s\n", lg.toString()));
}
logger.info("Serializing media library...");
mediaLibrary.serialize(oos);
logoutput.append("Serialized " + mediaLibrary.allTracks().size() + " media tracks\n");
for (MediaTrack lg : mediaLibrary.allTracks()) {
logoutput.append(String.format("%3s\n", lg.toString()));
}
logger.info(logoutput.toString());
}
@SuppressWarnings("unchecked") // the type-parameters in the casts of the de-serialized collection objects can't be checked
@Override
public void initiallyFillFrom(ObjectInputStream ois) throws IOException, ClassNotFoundException, InterruptedException {
logger.info("Performing initial replication load on "+this);
logger.info("Performing initial replication load on " + this);
ClassLoader oldContextClassloader = Thread.currentThread().getContextClassLoader();
try {
// Use this object's class's class loader as the context class loader which will then be used for
@@ -1674,9 +1704,11 @@ public class RacingEventServiceImpl implements RacingEventService, RegattaListen
regattasByName.clear();
regattasObservedForDefaultLeaderboard.clear();
for (DynamicTrackedRegatta regatta: regattaTrackingCache.values()) {
for (RaceTracker tracker : raceTrackersByRegatta.get(regatta)) {
tracker.stop();
if (raceTrackersByRegatta != null && !raceTrackersByRegatta.isEmpty()) {
for (DynamicTrackedRegatta regatta: regattaTrackingCache.values()) {
for (RaceTracker tracker : raceTrackersByRegatta.get(regatta)) {
tracker.stop();
}
}
}
@@ -1685,40 +1717,62 @@ public class RacingEventServiceImpl implements RacingEventService, RegattaListen
leaderboardsByName.clear();
eventsById.clear();
mediaLibrary.clear();
logger.info("receiving eventsById");
StringBuffer logoutput = new StringBuffer();
eventsById.putAll((Map<Serializable, Event>) ois.readObject());
logger.info("Recieved " + eventsById.size() + " NEW events");
logger.info("receiving regattasByName");
logoutput.append("\nReceived " + eventsById.size() + " NEW events\n");
for (Event event : eventsById.values()) {
logoutput.append(String.format("%3s\n", event.toString()));
}
regattasByName.putAll((Map<String, Regatta>) ois.readObject());
logger.info("Recieved " + regattasByName.size() + " NEW regattas");
logoutput.append("Received " + regattasByName.size() + " NEW regattas\n");
for (Regatta regatta : regattasByName.values()) {
logoutput.append(String.format("%3s\n", regatta.toString()));
}
// it is important that the leaderboards and tracked regattas are cleared before auto-linking to
// old leaderboards takes place which then don't match the new ones
for (DynamicTrackedRegatta trackedRegattaToObserve : (Set<DynamicTrackedRegatta>) ois.readObject()) {
ensureRegattaIsObservedForDefaultLeaderboardAndAutoLeaderboardLinking(trackedRegattaToObserve);
}
logger.info("receiving regattaTrackingCache");
regattaTrackingCache.putAll((Map<Regatta, DynamicTrackedRegatta>) ois.readObject());
logger.info("Recieved " + regattaTrackingCache.size() + " NEW regatta tracking cache entries");
logger.info("receiving leaderboardGroupsByName");
logoutput.append("Received " + regattaTrackingCache.size() + " NEW regatta tracking cache entries\n");
leaderboardGroupsByName.putAll((Map<String, LeaderboardGroup>) ois.readObject());
logger.info("Received " + leaderboardGroupsByName.size() + " NEW leaderboard groups");
logger.info("receiving leaderboardsByName");
logoutput.append("Received " + leaderboardGroupsByName.size() + " NEW leaderboard groups\n");
for (LeaderboardGroup lg : leaderboardGroupsByName.values()) {
logoutput.append(String.format("%3s\n", lg.toString()));
}
leaderboardsByName.putAll((Map<String, Leaderboard>) ois.readObject());
logger.info("Recieved " + leaderboardsByName.size() + " NEW leaderboards");
logoutput.append("Received " + leaderboardsByName.size() + " NEW leaderboards\n");
for (Leaderboard leaderboard : leaderboardsByName.values()) {
logoutput.append(String.format("%3s\n", leaderboard.toString()));
}
// now fix ScoreCorrectionListener setup for LeaderboardGroupMetaLeaderboard instances:
for (Leaderboard leaderboard : leaderboardsByName.values()) {
if (leaderboard instanceof LeaderboardGroupMetaLeaderboard) {
((LeaderboardGroupMetaLeaderboard) leaderboard).registerAsScoreCorrectionChangeForwarderAndRaceColumnListenerOnAllLeaderboards();
}
}
logger.info("receiving mediaLibrary");
mediaLibrary.deserialize(ois);
logger.info("Done with initial replication on "+this);
logoutput.append("Received " + mediaLibrary.allTracks().size() + " NEW media tracks\n");
for (MediaTrack mediatrack : mediaLibrary.allTracks()) {
logoutput.append(String.format("%3s\n", mediatrack.toString()));
}
logger.info(logoutput.toString());
logger.info("Done with initial replication on " + this);
} finally {
Thread.currentThread().setContextClassLoader(oldContextClassloader);
}
}
@Override
public long getDelayToLiveInMillis() {
return delayToLiveInMillis;
@@ -1886,5 +1940,9 @@ public class RacingEventServiceImpl implements RacingEventService, RegattaListen
public Collection<MediaTrack> getAllMediaTracks() {
return mediaLibrary.allTracks();
}
public String toString() {
return "RacingEventService: " + this.hashCode() + " Build: " + BuildVersion.getBuildVersion();
}
}
@@ -25,9 +25,11 @@ public abstract class AbstractRecordRaceLogEvent extends AbstractRacingEventServ
* @return <code>true</code> if add was successful.
*/
protected RaceLogEvent addEventTo(RaceColumn raceColumn) {
Fleet fleet = raceColumn.getFleetByName(fleetName);
RaceLog raceLog = raceColumn.getRaceLog(fleet);
raceLog.add(event);
if (raceColumn != null) {
Fleet fleet = raceColumn.getFleetByName(fleetName);
RaceLog raceLog = raceColumn.getRaceLog(fleet);
raceLog.add(event);
}
return event;
}
+2 -2
View File
@@ -90,10 +90,10 @@ a {
min-height: 100%; }
#yield_content .headNavigation {
float: right;
margin: 20px 0 0 0;
margin: -30px 0 0 0;
padding: 0; }
#yield_content .headNavigation a {
font-family: 'OpenSans-Bold';
font-family: 'OpenSans-Bold', Arial, Verdana, sans-serif;
color: #fff; }
/* line 47, site.css.sass */
#yield_content #pageHeadline {
+1 -1
View File
@@ -293,7 +293,7 @@
<repository>
<id>sailing-server-maven</id>
<url>http://sapcoe-app01.pironet-ndh.com/maven</url>
<url>http://maven.sapsailing.com/maven</url>
</repository>
</repositories>
</project>