From 0601becf8185982a48f39358758a5cd9542fa3ae Mon Sep 17 00:00:00 2001 From: Steffen Schaefer Date: Wed, 21 Mar 2018 16:50:37 +0100 Subject: [PATCH] Bug 4489: Initial implementation to ensure replication state for out of thread event processing in TrackedRegattaImpl --- .../impl/CreateAndTrackWithRaceLogTest.java | 3 +- .../impl/RaceLogRaceTracker.java | 4 +- .../impl/SwissTimingRaceTrackerImpl.java | 4 +- .../SwissTimingReplayToDomainAdapter.java | 4 +- .../test/FetchTracksAndStoreLocallyTest.java | 3 +- .../domain/test/ReceiveTrackingDataTest.java | 3 +- .../test/TestDeadlockInRegattaListener.java | 11 ++-- .../domain/test/mock/MockedTrackedRace.java | 10 ++-- .../tracking/impl/TrackedRegattaTest.java | 7 ++- .../impl/DomainFactoryImpl.java | 3 +- .../impl/RaceCourseReceiver.java | 4 +- .../tractracadapter/impl/Simulator.java | 3 +- ...AbstractTrackedRegattaAndRaceObserver.java | 3 +- .../domain/tracking/TrackedRegatta.java | 11 ++-- .../impl/DynamicTrackedRegattaImpl.java | 8 ++- .../domain/tracking/impl/TrackedRaceImpl.java | 2 +- .../tracking/impl/TrackedRegattaImpl.java | 56 ++++++++++++------- .../gwt/ui/test/MockedTrackedRace.java | 10 ++-- ...estStoringAndRetrievingWindTracksTest.java | 4 +- .../server/statistics/StatisticsTest.java | 5 +- .../server/test/RaceTrackerStartStopTest.java | 10 +++- .../sailing/server/test/RaceTrackerTest.java | 3 +- .../server/test/RemoveLeaderboardTest.java | 4 +- .../server/impl/RacingEventServiceImpl.java | 10 +++- .../ImportMasterDataOperation.java | 4 +- .../simulation/SimulationServiceImpl.java | 3 +- 26 files changed, 128 insertions(+), 64 deletions(-) diff --git a/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/CreateAndTrackWithRaceLogTest.java b/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/CreateAndTrackWithRaceLogTest.java index 0e20c9b4525..21a9059b81e 100644 --- a/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/CreateAndTrackWithRaceLogTest.java +++ b/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/CreateAndTrackWithRaceLogTest.java @@ -11,6 +11,7 @@ import java.io.IOException; import java.net.MalformedURLException; import java.net.URISyntaxException; import java.util.Collections; +import java.util.Optional; import java.util.UUID; import org.junit.After; @@ -228,7 +229,7 @@ public class CreateAndTrackWithRaceLogTest extends RaceLogTrackingTestHelper { public void raceAdded(TrackedRace trackedRace) { } }; - raceHandle.getTrackedRegatta().addRaceListener(raceListener); + raceHandle.getTrackedRegatta().addRaceListener(raceListener, Optional.empty()); raceHandle.getTrackedRegatta().removeRaceListener(raceListener).get(); } diff --git a/java/com.sap.sailing.domain.racelogtrackingadapter/src/com/sap/sailing/domain/racelogtracking/impl/RaceLogRaceTracker.java b/java/com.sap.sailing.domain.racelogtrackingadapter/src/com/sap/sailing/domain/racelogtracking/impl/RaceLogRaceTracker.java index 395d6a2e21a..e354df7ba7e 100755 --- a/java/com.sap.sailing.domain.racelogtrackingadapter/src/com/sap/sailing/domain/racelogtracking/impl/RaceLogRaceTracker.java +++ b/java/com.sap.sailing.domain.racelogtrackingadapter/src/com/sap/sailing/domain/racelogtracking/impl/RaceLogRaceTracker.java @@ -7,6 +7,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.Map.Entry; import java.util.logging.Level; import java.util.logging.Logger; @@ -304,7 +305,8 @@ public class RaceLogRaceTracker extends AbstractRaceTrackerBaseImpl { raceColumn.setRaceIdentifier(fleet, trackedRegatta.getRegatta().getRaceIdentifier(raceDef)); trackedRace = trackedRegatta.createTrackedRace(raceDef, sidelines, windStore, params.getDelayToLiveInMillis(), WindTrack.DEFAULT_MILLISECONDS_OVER_WHICH_TO_AVERAGE_WIND, - boatClass.getApproximateManeuverDurationInMilliseconds(), null, /*useMarkPassingCalculator*/ true, raceLogResolver); + boatClass.getApproximateManeuverDurationInMilliseconds(), null, /*useMarkPassingCalculator*/ true, raceLogResolver, + /* Not needed because the RaceTracker is not active on a replica */ Optional.empty()); notifyRaceCreationListeners(); logger.info(String.format("Started tracking race-log race (%s)", raceLog)); // this wakes up all waiting race handles diff --git a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingRaceTrackerImpl.java b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingRaceTrackerImpl.java index 18bd0aee11a..87881b9e7b1 100644 --- a/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingRaceTrackerImpl.java +++ b/java/com.sap.sailing.domain.swisstimingadapter/src/com/sap/sailing/domain/swisstimingadapter/impl/SwissTimingRaceTrackerImpl.java @@ -10,6 +10,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.NavigableSet; +import java.util.Optional; import java.util.TreeMap; import java.util.logging.Level; import java.util.logging.Logger; @@ -487,7 +488,8 @@ public class SwissTimingRaceTrackerImpl extends AbstractRaceTrackerImpl // we already know our single RaceDefinition assert SwissTimingRaceTrackerImpl.this.race == race; } - }, useInternalMarkPassingAlgorithm, raceLogResolver); + }, useInternalMarkPassingAlgorithm, raceLogResolver, + /* Not needed because the RaceTracker is not active on a replica */ Optional.empty()); notifyRaceCreationListeners(); logger.info("Created SwissTiming RaceDefinition and TrackedRace for "+race.getName()); } diff --git a/java/com.sap.sailing.domain.swisstimingreplayadapter/src/com/sap/sailing/domain/swisstimingreplayadapter/impl/SwissTimingReplayToDomainAdapter.java b/java/com.sap.sailing.domain.swisstimingreplayadapter/src/com/sap/sailing/domain/swisstimingreplayadapter/impl/SwissTimingReplayToDomainAdapter.java index aa3489242a9..dee036575db 100755 --- a/java/com.sap.sailing.domain.swisstimingreplayadapter/src/com/sap/sailing/domain/swisstimingreplayadapter/impl/SwissTimingReplayToDomainAdapter.java +++ b/java/com.sap.sailing.domain.swisstimingreplayadapter/src/com/sap/sailing/domain/swisstimingreplayadapter/impl/SwissTimingReplayToDomainAdapter.java @@ -11,6 +11,7 @@ import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.NavigableSet; +import java.util.Optional; import java.util.Set; import java.util.TimeZone; import java.util.logging.Logger; @@ -400,7 +401,8 @@ public class SwissTimingReplayToDomainAdapter extends SwissTimingReplayAdapter i TrackedRace.DEFAULT_LIVE_DELAY_IN_MILLISECONDS, WindTrack.DEFAULT_MILLISECONDS_OVER_WHICH_TO_AVERAGE_WIND, /* time over which to average speed: */ race.getBoatClass().getApproximateManeuverDurationInMilliseconds(), - /* raceDefinitionSetToUpdate */ null, useInternalMarkPassingAlgorithm, raceLogResolver); + /* raceDefinitionSetToUpdate */ null, useInternalMarkPassingAlgorithm, raceLogResolver, + /* Not needed because the RaceTracker is not active on a replica */ Optional.empty()); trackedRace.onStatusChanged(this, new TrackedRaceStatusImpl(TrackedRaceStatusEnum.LOADING, 0)); TimePoint bestStartTimeKnownSoFar = bestStartTimePerRaceID.get(currentRaceID); if (bestStartTimeKnownSoFar != null) { diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/FetchTracksAndStoreLocallyTest.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/FetchTracksAndStoreLocallyTest.java index d965b5fcc1a..7b8fa85d92c 100755 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/FetchTracksAndStoreLocallyTest.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/FetchTracksAndStoreLocallyTest.java @@ -6,6 +6,7 @@ import java.net.MalformedURLException; import java.net.URISyntaxException; import java.util.HashMap; import java.util.Map; +import java.util.Optional; import org.junit.Ignore; import org.junit.Test; @@ -74,7 +75,7 @@ public class FetchTracksAndStoreLocallyTest extends OnlineTracTracBasedTest { @Override public void raceRemoved(TrackedRace trackedRace) { } - }); + }, Optional.empty()); super.completeSetupLaunchingControllerAndWaitForRaceDefinition(ReceiverType.RACECOURSE, ReceiverType.RACESTARTFINISH, ReceiverType.RAWPOSITIONS); } diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/ReceiveTrackingDataTest.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/ReceiveTrackingDataTest.java index 11d389f28a1..855b29f764a 100755 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/ReceiveTrackingDataTest.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/ReceiveTrackingDataTest.java @@ -5,6 +5,7 @@ import static org.mockito.Mockito.mock; import java.net.MalformedURLException; import java.net.URISyntaxException; +import java.util.Optional; import org.junit.Before; import org.junit.Test; @@ -77,7 +78,7 @@ public class ReceiveTrackingDataTest extends AbstractTracTracLiveTest { @Override public void raceRemoved(TrackedRace trackedRace) { } - }); + }, Optional.empty()); for (Receiver receiver : domainFactory .getUpdateReceivers(trackedRegatta, /* delayToLiveInMillis */0l, /* simulator */null, EmptyWindStore.INSTANCE, new DynamicRaceDefinitionSet() { diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/TestDeadlockInRegattaListener.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/TestDeadlockInRegattaListener.java index 54b0889169d..c7020e2d91e 100755 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/TestDeadlockInRegattaListener.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/TestDeadlockInRegattaListener.java @@ -5,6 +5,7 @@ import static org.mockito.Mockito.when; import java.io.IOException; import java.net.MalformedURLException; +import java.util.Optional; import java.util.UUID; import java.util.concurrent.BrokenBarrierException; import java.util.concurrent.CyclicBarrier; @@ -33,6 +34,7 @@ import com.sap.sailing.domain.tracking.impl.DynamicTrackedRaceImpl; import com.sap.sailing.domain.tracking.impl.DynamicTrackedRegattaImpl; import com.sap.sailing.server.RacingEventService; import com.sap.sailing.server.impl.RacingEventServiceImpl; +import com.sap.sse.util.ThreadLocalTransporter; public class TestDeadlockInRegattaListener { @Rule @@ -56,13 +58,14 @@ public class TestDeadlockInRegattaListener { private static final long serialVersionUID = -3599667964201700780L; @Override - protected void notifyListenersAboutTrackedRaceRemoved(TrackedRace trackedRace) { + protected void notifyListenersAboutTrackedRaceRemoved(TrackedRace trackedRace, + Optional threadLocalTransporter) { try { latch.await(); } catch (InterruptedException | BrokenBarrierException e) { throw new RuntimeException(e); } - super.notifyListenersAboutTrackedRaceRemoved(trackedRace); + super.notifyListenersAboutTrackedRaceRemoved(trackedRace, Optional.empty()); } }; RacingEventServiceImpl racingEventService = new RacingEventServiceImpl() { @@ -125,10 +128,10 @@ public class TestDeadlockInRegattaListener { throw new RuntimeException(e); } }).start(); - trackedRegatta.addTrackedRace(trackedRace1); + trackedRegatta.addTrackedRace(trackedRace1, Optional.empty()); // the following runs into RacingEventService.getRaceTrackerByRegattaAndRaceIdentifier // which waits for the latch based on the override above while in synchronized RegattaListener.raceAdded - new Thread(()->trackedRegatta.addTrackedRace(trackedRace2)).start(); + new Thread(()->trackedRegatta.addTrackedRace(trackedRace2, Optional.empty())).start(); monitorOnRegattaListenerLatch.await(); // the following awaits the latch in TrackedRegattaImpl.notifyListenersAboutTrackedRaceRemoved // after the write lock has been obtained but before the synchronized RegattaListener.raceRemoved method diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/mock/MockedTrackedRace.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/mock/MockedTrackedRace.java index e4a45f5d83b..578c9aec53a 100644 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/mock/MockedTrackedRace.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/mock/MockedTrackedRace.java @@ -4,6 +4,7 @@ import java.io.Serializable; import java.util.Collections; import java.util.List; import java.util.NavigableSet; +import java.util.Optional; import java.util.Set; import java.util.TreeSet; import java.util.concurrent.Future; @@ -89,6 +90,7 @@ import com.sap.sse.common.IsManagedByCache; import com.sap.sse.common.TimePoint; import com.sap.sse.common.Util; import com.sap.sse.common.Util.Pair; +import com.sap.sse.util.ThreadLocalTransporter; public class MockedTrackedRace implements DynamicTrackedRace { private static final long serialVersionUID = 5827912985564121181L; @@ -565,15 +567,15 @@ public class MockedTrackedRace implements DynamicTrackedRace { } @Override - public void addTrackedRace(TrackedRace trackedRace) { + public void addTrackedRace(TrackedRace trackedRace, Optional threadLocalTransporter) { } @Override - public void removeTrackedRace(TrackedRace trackedRace) { + public void removeTrackedRace(TrackedRace trackedRace, Optional threadLocalTransporter) { } @Override - public void addRaceListener(RaceListener listener) { + public void addRaceListener(RaceListener listener, Optional threadLocalTransporter) { } @Override @@ -596,7 +598,7 @@ public class MockedTrackedRace implements DynamicTrackedRace { WindStore windStore, long delayToLiveInMillis, long millisecondsOverWhichToAverageWind, long millisecondsOverWhichToAverageSpeed, DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useMarkPassingcalculator, - RaceLogResolver raceLogResolver) { + RaceLogResolver raceLogResolver, Optional threadLocalTransporter) { return null; } diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/tracking/impl/TrackedRegattaTest.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/tracking/impl/TrackedRegattaTest.java index 0c18176671a..3980819bac9 100644 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/tracking/impl/TrackedRegattaTest.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/tracking/impl/TrackedRegattaTest.java @@ -4,6 +4,7 @@ import static org.mockito.Mockito.mock; import java.util.Arrays; import java.util.Collections; +import java.util.Optional; import java.util.concurrent.CyclicBarrier; import java.util.concurrent.Phaser; import java.util.concurrent.TimeUnit; @@ -74,11 +75,11 @@ public class TrackedRegattaTest { throw new RuntimeException(e); } } - }); + }, Optional.empty()); DynamicTrackedRace race1 = createRace("R1"); Thread thread1 = new Thread(() -> { - regatta.addTrackedRace(race1); + regatta.addTrackedRace(race1, Optional.empty()); }); thread1.start(); // This ensures, that the add event is being processed but is not finished because @@ -88,7 +89,7 @@ public class TrackedRegattaTest { addPhaser.arriveAndAwaitAdvance(); Thread thread2 = new Thread(() -> { - regatta.removeTrackedRace(race1); + regatta.removeTrackedRace(race1, Optional.empty()); }); thread2.start(); // If the implementation ensures that the events are fired in order, diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/DomainFactoryImpl.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/DomainFactoryImpl.java index c65b8bd83cd..a2cc640340e 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/DomainFactoryImpl.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/DomainFactoryImpl.java @@ -13,6 +13,7 @@ import java.util.HashSet; import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.Map.Entry; import java.util.Set; import java.util.UUID; @@ -616,7 +617,7 @@ public class DomainFactoryImpl implements DomainFactory { return trackedRegatta.createTrackedRace(race, sidelines, windStore, delayToLiveInMillis, millisecondsOverWhichToAverageWind, /* time over which to average speed: */ race.getBoatClass().getApproximateManeuverDurationInMilliseconds(), - raceDefinitionSetToUpdate, useMarkPassingCalculator, raceLogResolver); + raceDefinitionSetToUpdate, useMarkPassingCalculator, raceLogResolver, Optional.empty()); } @Override diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/RaceCourseReceiver.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/RaceCourseReceiver.java index 8487db75d19..1ae14c249a3 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/RaceCourseReceiver.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/RaceCourseReceiver.java @@ -5,6 +5,7 @@ import java.util.ArrayList; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.function.Consumer; import java.util.logging.Level; import java.util.logging.Logger; @@ -241,7 +242,8 @@ public class RaceCourseReceiver extends AbstractReceiverWithQueue sidelines, WindStore windStore, long delayToLiveInMillis, long millisecondsOverWhichToAverageWind, long millisecondsOverWhichToAverageSpeed, - DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useInternalMarkPassingAlgorithm, RaceLogResolver raceLogResolver); + DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useInternalMarkPassingAlgorithm, RaceLogResolver raceLogResolver, + Optional beforeAndAfterNotificationHandler); /** * Obtains the tracked race for race. Blocks until the tracked race has been created @@ -79,16 +82,16 @@ public interface TrackedRegatta extends Serializable { */ TrackedRace getExistingTrackedRace(RaceDefinition race); - void addTrackedRace(TrackedRace trackedRace); + void addTrackedRace(TrackedRace trackedRace, Optional beforeAndAfterNotificationHandler); - void removeTrackedRace(TrackedRace trackedRace); + void removeTrackedRace(TrackedRace trackedRace, Optional beforeAndAfterNotificationHandler); /** * Listener will be notified when {@link #addTrackedRace(TrackedRace)} is called and * upon registration for each tracked race already known. Therefore, the listener * won't miss any tracked race. */ - void addRaceListener(RaceListener listener); + void addRaceListener(RaceListener listener, Optional beforeAndAfterNotificationHandler); /** * Removes the given listener and returns a {@link Future} that will be completed diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/DynamicTrackedRegattaImpl.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/DynamicTrackedRegattaImpl.java index fdc820b7904..990fa8cb327 100755 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/DynamicTrackedRegattaImpl.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/DynamicTrackedRegattaImpl.java @@ -1,5 +1,7 @@ package com.sap.sailing.domain.tracking.impl; +import java.util.Optional; + import com.sap.sailing.domain.abstractlog.race.analyzing.impl.RaceLogResolver; import com.sap.sailing.domain.base.RaceDefinition; import com.sap.sailing.domain.base.Regatta; @@ -8,6 +10,7 @@ import com.sap.sailing.domain.tracking.DynamicRaceDefinitionSet; import com.sap.sailing.domain.tracking.DynamicTrackedRace; import com.sap.sailing.domain.tracking.DynamicTrackedRegatta; import com.sap.sailing.domain.tracking.WindStore; +import com.sap.sse.util.ThreadLocalTransporter; public class DynamicTrackedRegattaImpl extends TrackedRegattaImpl implements DynamicTrackedRegatta { private static final long serialVersionUID = -90155868534737120L; @@ -35,9 +38,10 @@ public class DynamicTrackedRegattaImpl extends TrackedRegattaImpl implements Dyn @Override public DynamicTrackedRace createTrackedRace(RaceDefinition raceDefinition, Iterable sidelines, WindStore windStore, long delayToLiveInMillis, long millisecondsOverWhichToAverageWind, long millisecondsOverWhichToAverageSpeed, - DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useMarkPassingCalculator, RaceLogResolver raceLogResolver) { + DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useMarkPassingCalculator, RaceLogResolver raceLogResolver, + Optional threadLocalTransporter) { return (DynamicTrackedRace) super.createTrackedRace(raceDefinition, sidelines, windStore, delayToLiveInMillis, millisecondsOverWhichToAverageWind, - millisecondsOverWhichToAverageSpeed, raceDefinitionSetToUpdate, useMarkPassingCalculator, raceLogResolver); + millisecondsOverWhichToAverageSpeed, raceDefinitionSetToUpdate, useMarkPassingCalculator, raceLogResolver, threadLocalTransporter); } } diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedRaceImpl.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedRaceImpl.java index 38ffed188a4..d82fc8e09d1 100755 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedRaceImpl.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedRaceImpl.java @@ -548,7 +548,7 @@ public abstract class TrackedRaceImpl extends TrackedRaceWithWindEssentials impl markPassingCalculator.stop(); } } - }); + }, /* Not relevant For replication */ Optional.empty()); } else { markPassingCalculator = null; } diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedRegattaImpl.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedRegattaImpl.java index e1185b9adfc..8585cdad67f 100755 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedRegattaImpl.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedRegattaImpl.java @@ -9,6 +9,7 @@ import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; @@ -34,6 +35,7 @@ import com.sap.sse.common.TimePoint; import com.sap.sse.common.Util; import com.sap.sse.concurrent.LockUtil; import com.sap.sse.concurrent.NamedReentrantReadWriteLock; +import com.sap.sse.util.ThreadLocalTransporter; public class TrackedRegattaImpl implements TrackedRegatta { private static final long serialVersionUID = 6480508193567014285L; @@ -128,7 +130,7 @@ public class TrackedRegattaImpl implements TrackedRegatta { } @Override - public void addTrackedRace(TrackedRace trackedRace) { + public void addTrackedRace(TrackedRace trackedRace, Optional threadLocalTransporter) { final TrackedRace oldTrackedRace; lockTrackedRacesForWrite(); try { @@ -136,39 +138,51 @@ public class TrackedRegattaImpl implements TrackedRegatta { " with regatta hash code "+getRegatta().hashCode()); oldTrackedRace = trackedRaces.put(trackedRace.getRace(), trackedRace); if (oldTrackedRace != trackedRace) { - notifyListenersAboutTrackedRaceAdded(trackedRace); + notifyListenersAboutTrackedRaceAdded(trackedRace, threadLocalTransporter); } } finally { unlockTrackedRacesAfterWrite(); } } - protected void notifyListenersAboutTrackedRaceAdded(TrackedRace trackedRace) { - enqueEvent(listener -> listener.raceAdded(trackedRace)); + protected void notifyListenersAboutTrackedRaceAdded(TrackedRace trackedRace, Optional threadLocalTransporter) { + enqueEvent(listener -> listener.raceAdded(trackedRace), threadLocalTransporter); } - protected void enqueEvent(Consumer fireEventCallback) { + protected void enqueEvent(Consumer fireEventCallback, Optional threadLocalTransporter) { final Set listenersToInform = new HashSet<>(raceListeners.keySet()); + threadLocalTransporter.ifPresent(ThreadLocalTransporter::rememberThreadLocalStates); eventQueue.addWork(() -> { - for (RaceListener listener : listenersToInform) { - fireEventCallback.accept(listener); - } + withBeforeAndAfterHandling(threadLocalTransporter, () -> { + for (RaceListener listener : listenersToInform) { + fireEventCallback.accept(listener); + } + }); }); } + private void withBeforeAndAfterHandling(Optional threadLocalTransporter, Runnable action) { + threadLocalTransporter.ifPresent(ThreadLocalTransporter::pushThreadLocalStates); + try { + action.run(); + } finally { + threadLocalTransporter.ifPresent(ThreadLocalTransporter::popThreadLocalStates); + } + } + @Override - public void removeTrackedRace(TrackedRace trackedRace) { + public void removeTrackedRace(TrackedRace trackedRace, Optional threadLocalTransporter) { lockTrackedRacesForWrite(); try { trackedRaces.remove(trackedRace.getRace()); - notifyListenersAboutTrackedRaceRemoved(trackedRace); + notifyListenersAboutTrackedRaceRemoved(trackedRace, threadLocalTransporter); } finally { unlockTrackedRacesAfterWrite(); } } - protected void notifyListenersAboutTrackedRaceRemoved(TrackedRace trackedRace) { - enqueEvent(listener -> listener.raceRemoved(trackedRace)); + protected void notifyListenersAboutTrackedRaceRemoved(TrackedRace trackedRace, Optional threadLocalTransporter) { + enqueEvent(listener -> listener.raceRemoved(trackedRace), threadLocalTransporter); } @Override @@ -201,7 +215,7 @@ public class TrackedRegattaImpl implements TrackedRegatta { } } }; - addRaceListener(listener); + addRaceListener(listener, Optional.empty()); try { synchronized (mutex) { result = getExistingTrackedRace(race); @@ -232,16 +246,19 @@ public class TrackedRegattaImpl implements TrackedRegatta { } @Override - public void addRaceListener(RaceListener listener) { + public void addRaceListener(RaceListener listener, Optional threadLocalTransporter) { lockTrackedRacesForRead(); try { raceListeners.put(listener, listener); final List trackedRacesCopy = new ArrayList<>(); Util.addAll(getTrackedRaces(), trackedRacesCopy); + threadLocalTransporter.ifPresent(ThreadLocalTransporter::rememberThreadLocalStates); eventQueue.addWork(() -> { - for (TrackedRace trackedRace : trackedRacesCopy) { - listener.raceAdded(trackedRace); - } + withBeforeAndAfterHandling(threadLocalTransporter, () -> { + for (TrackedRace trackedRace : trackedRacesCopy) { + listener.raceAdded(trackedRace); + } + }); }); } finally { unlockTrackedRacesAfterRead(); @@ -284,7 +301,8 @@ public class TrackedRegattaImpl implements TrackedRegatta { public DynamicTrackedRace createTrackedRace(RaceDefinition raceDefinition, Iterable sidelines, WindStore windStore, long delayToLiveInMillis, long millisecondsOverWhichToAverageWind, long millisecondsOverWhichToAverageSpeed, - DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useInternalMarkPassingAlgorithm, RaceLogResolver raceLogResolver) { + DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useInternalMarkPassingAlgorithm, RaceLogResolver raceLogResolver, + Optional threadLocalTransporter) { logger.log(Level.INFO, "Creating DynamicTrackedRaceImpl for RaceDefinition " + raceDefinition.getName()); DynamicTrackedRaceImpl result = new DynamicTrackedRaceImpl(this, raceDefinition, sidelines, windStore, delayToLiveInMillis, millisecondsOverWhichToAverageWind, @@ -295,7 +313,7 @@ public class TrackedRegattaImpl implements TrackedRegatta { if (raceDefinitionSetToUpdate != null) { raceDefinitionSetToUpdate.addRaceDefinition(raceDefinition, result); } - addTrackedRace(result); + addTrackedRace(result, threadLocalTransporter); return result; } } diff --git a/java/com.sap.sailing.gwt.ui.test/src/com/sap/sailing/gwt/ui/test/MockedTrackedRace.java b/java/com.sap.sailing.gwt.ui.test/src/com/sap/sailing/gwt/ui/test/MockedTrackedRace.java index 201f21b1cc5..d0b893acc81 100644 --- a/java/com.sap.sailing.gwt.ui.test/src/com/sap/sailing/gwt/ui/test/MockedTrackedRace.java +++ b/java/com.sap.sailing.gwt.ui.test/src/com/sap/sailing/gwt/ui/test/MockedTrackedRace.java @@ -7,6 +7,7 @@ import java.io.Serializable; import java.util.Collections; import java.util.List; import java.util.NavigableSet; +import java.util.Optional; import java.util.Set; import java.util.concurrent.Future; @@ -76,6 +77,7 @@ import com.sap.sse.common.Duration; import com.sap.sse.common.IsManagedByCache; import com.sap.sse.common.TimePoint; import com.sap.sse.common.Util; +import com.sap.sse.util.ThreadLocalTransporter; public class MockedTrackedRace implements DynamicTrackedRace { private static final long serialVersionUID = 5827912985564121181L; @@ -295,15 +297,15 @@ public class MockedTrackedRace implements DynamicTrackedRace { } @Override - public void addTrackedRace(TrackedRace trackedRace) { + public void addTrackedRace(TrackedRace trackedRace, Optional threadLocalTransporter) { } @Override - public void removeTrackedRace(TrackedRace trackedRace) { + public void removeTrackedRace(TrackedRace trackedRace, Optional threadLocalTransporter) { } @Override - public void addRaceListener(RaceListener listener) { + public void addRaceListener(RaceListener listener, Optional threadLocalTransporter) { } @Override @@ -325,7 +327,7 @@ public class MockedTrackedRace implements DynamicTrackedRace { public DynamicTrackedRace createTrackedRace(RaceDefinition raceDefinition, Iterable sidelines, WindStore windStore, long delayToLiveInMillis, long millisecondsOverWhichToAverageWind, long millisecondsOverWhichToAverageSpeed, DynamicRaceDefinitionSet raceDefinitionSetToUpdate, - boolean useMarkPassingCalculator, RaceLogResolver raceLogResolver) { + boolean useMarkPassingCalculator, RaceLogResolver raceLogResolver, Optional threadLocalTransporter) { return null; } diff --git a/java/com.sap.sailing.mongodb.test/src/com/sap/sailing/mongodb/test/TestStoringAndRetrievingWindTracksTest.java b/java/com.sap.sailing.mongodb.test/src/com/sap/sailing/mongodb/test/TestStoringAndRetrievingWindTracksTest.java index f78772ff244..e53c92b45da 100755 --- a/java/com.sap.sailing.mongodb.test/src/com/sap/sailing/mongodb/test/TestStoringAndRetrievingWindTracksTest.java +++ b/java/com.sap.sailing.mongodb.test/src/com/sap/sailing/mongodb/test/TestStoringAndRetrievingWindTracksTest.java @@ -9,6 +9,7 @@ import java.net.MalformedURLException; import java.net.URISyntaxException; import java.net.UnknownHostException; import java.util.Collections; +import java.util.Optional; import org.junit.Before; import org.junit.Test; @@ -103,7 +104,8 @@ public class TestStoringAndRetrievingWindTracksTest extends AbstractTracTracLive @Override public void addRaceDefinition(RaceDefinition race, DynamicTrackedRace trackedRace) { } - }, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class)); + }, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class), + Optional.empty()); WindSource windSource = new WindSourceImpl(WindSourceType.WEB); Mongo myFirstMongo = newMongo(); DB firstDatabase = myFirstMongo.getDB(dbConfiguration.getDatabaseName()); diff --git a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/statistics/StatisticsTest.java b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/statistics/StatisticsTest.java index 46bbae978ea..17e91a94394 100644 --- a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/statistics/StatisticsTest.java +++ b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/statistics/StatisticsTest.java @@ -8,6 +8,7 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.List; +import java.util.Optional; import org.junit.Before; import org.junit.Test; @@ -81,7 +82,7 @@ public class StatisticsTest { trackedRace.setEndOfTrackingReceived(new MillisecondsTimePoint(END_OF_TRACKING)); trackedRace.setStartTimeReceived(new MillisecondsTimePoint(START_OF_RACE)); - regatta.addTrackedRace(trackedRace); + regatta.addTrackedRace(trackedRace, Optional.empty()); } private TrackedRaceStatisticsCacheImpl getStatisticsCacheWithRegattaAdded() throws Exception { @@ -98,7 +99,7 @@ public class StatisticsTest { public void raceAdded(TrackedRace trackedRace) { } }; - regatta.addRaceListener(raceListener); + regatta.addRaceListener(raceListener, Optional.empty()); regatta.removeRaceListener(raceListener).get(); return trackedRaceStatisticsCache; diff --git a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RaceTrackerStartStopTest.java b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RaceTrackerStartStopTest.java index 16f1d062661..8685d9ca7a5 100644 --- a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RaceTrackerStartStopTest.java +++ b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RaceTrackerStartStopTest.java @@ -14,6 +14,7 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.Iterator; +import java.util.Optional; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -83,17 +84,20 @@ public class RaceTrackerStartStopTest { trackedRegatta1.createTrackedRace(raceDef1, Collections. emptyList(), /* windStore */ EmptyWindStore.INSTANCE, /* delayToLiveInMillis */ 0l, /* millisecondsOverWhichToAverageWind */ 0l, - /* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class)); + /* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class), + Optional.empty()); regatta.addRace(raceDef2); trackedRegatta1.createTrackedRace(raceDef2, Collections. emptyList(), /* windStore */ EmptyWindStore.INSTANCE, /* delayToLiveInMillis */ 0l, /* millisecondsOverWhichToAverageWind */ 0l, - /* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class)); + /* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class), + Optional.empty()); regatta.addRace(raceDef3); trackedRegatta1.createTrackedRace(raceDef3, Collections. emptyList(), /* windStore */ EmptyWindStore.INSTANCE, /* delayToLiveInMillis */ 0l, /* millisecondsOverWhichToAverageWind */ 0l, - /* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class)); + /* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class), + Optional.empty()); Long trackerID1 = new Long(1); Long trackerID2 = new Long(2); Long trackerID3 = new Long(3); diff --git a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RaceTrackerTest.java b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RaceTrackerTest.java index b94b3655167..96509d108d4 100755 --- a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RaceTrackerTest.java +++ b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RaceTrackerTest.java @@ -10,6 +10,7 @@ import java.net.MalformedURLException; import java.net.URI; import java.net.URISyntaxException; import java.net.URL; +import java.util.Optional; import java.util.logging.Logger; import org.junit.After; @@ -102,7 +103,7 @@ public class RaceTrackerTest { @Override public void raceRemoved(TrackedRace trackedRace) { } - }); + }, Optional.empty()); synchronized (trackedRaces) { if (trackedRaces[0] == null) { trackedRaces.wait(); diff --git a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RemoveLeaderboardTest.java b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RemoveLeaderboardTest.java index 8a7d341470a..39a522db57f 100755 --- a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RemoveLeaderboardTest.java +++ b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/RemoveLeaderboardTest.java @@ -10,6 +10,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.UUID; import org.junit.Before; @@ -101,7 +102,8 @@ public class RemoveLeaderboardTest { trackedRace = trackedRegatta1.createTrackedRace(raceDef1, Collections. emptyList(), /* windStore */ EmptyWindStore.INSTANCE, /* delayToLiveInMillis */ 0l, /* millisecondsOverWhichToAverageWind */ 0l, - /* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class)); + /* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class), + Optional.empty()); } @Test diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/RacingEventServiceImpl.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/RacingEventServiceImpl.java index 284ea401e52..bb8a96c70e6 100644 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/RacingEventServiceImpl.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/RacingEventServiceImpl.java @@ -25,6 +25,7 @@ import java.util.List; import java.util.Locale; import java.util.Map; import java.util.Map.Entry; +import java.util.Optional; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -1740,13 +1741,15 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes RaceDefinition race = getRace(raceIdentifier); return trackedRegatta.createTrackedRace(race, Collections. emptyList(), windStore, delayToLiveInMillis, millisecondsOverWhichToAverageWind, millisecondsOverWhichToAverageSpeed, - /* raceDefinitionSetToUpdate */null, useMarkPassingCalculator, /* raceLogResolver */ this); + /* raceDefinitionSetToUpdate */null, useMarkPassingCalculator, /* raceLogResolver */ this, + Optional.of(this.getThreadLocalTransporterForCurrentlyFillingFromInitialLoadOrApplyingOperationReceivedFromMaster())); } private void ensureRegattaIsObservedForDefaultLeaderboardAndAutoLeaderboardLinking( DynamicTrackedRegatta trackedRegatta) { if (regattasObservedForDefaultLeaderboard.add(trackedRegatta)) { - trackedRegatta.addRaceListener(new RaceAdditionListener()); + trackedRegatta.addRaceListener(new RaceAdditionListener(), + Optional.of(this.getThreadLocalTransporterForCurrentlyFillingFromInitialLoadOrApplyingOperationReceivedFromMaster())); } } @@ -2436,7 +2439,8 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes final int newSizeOfTrackedRaces; oldSizeOfTrackedRaces = Util.size(trackedRegatta.getTrackedRaces()); try { - trackedRegatta.removeTrackedRace(trackedRace); + trackedRegatta.removeTrackedRace(trackedRace, Optional.of( + getThreadLocalTransporterForCurrentlyFillingFromInitialLoadOrApplyingOperationReceivedFromMaster())); newSizeOfTrackedRaces = Util.size(trackedRegatta.getTrackedRaces()); isTrackedRacesBecameEmpty = (oldSizeOfTrackedRaces > 0 && newSizeOfTrackedRaces == 0); } finally { diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/operationaltransformation/ImportMasterDataOperation.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/operationaltransformation/ImportMasterDataOperation.java index 7c24ea62cc0..a1983633ad7 100755 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/operationaltransformation/ImportMasterDataOperation.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/operationaltransformation/ImportMasterDataOperation.java @@ -7,6 +7,7 @@ import java.util.Collection; import java.util.List; import java.util.Map; import java.util.Map.Entry; +import java.util.Optional; import java.util.Set; import java.util.UUID; import java.util.logging.Level; @@ -477,7 +478,8 @@ public class ImportMasterDataOperation extends trackedRegatta.unlockTrackedRacesAfterRead(); } for (TrackedRace raceToRemove : toRemove) { - trackedRegatta.removeTrackedRace(raceToRemove); + trackedRegatta.removeTrackedRace(raceToRemove, Optional.of(toState + .getThreadLocalTransporterForCurrentlyFillingFromInitialLoadOrApplyingOperationReceivedFromMaster())); RaceDefinition race = existingRegatta.getRaceByName(raceToRemove .getRaceIdentifier().getRaceName()); if (race != null) { diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/simulation/SimulationServiceImpl.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/simulation/SimulationServiceImpl.java index f7ebb1eca4a..9252bbb57c4 100644 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/simulation/SimulationServiceImpl.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/simulation/SimulationServiceImpl.java @@ -6,6 +6,7 @@ import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Map.Entry; +import java.util.Optional; import java.util.TreeMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; @@ -283,7 +284,7 @@ public class SimulationServiceImpl implements SimulationService { DynamicTrackedRegatta trackedRegatta = racingEventService.getTrackedRegatta(regatta); SimulationRaceListener raceListener = new SimulationRaceListener(); raceListeners.put(legIdentifier.getRegattaName(), raceListener); - trackedRegatta.addRaceListener(raceListener); + trackedRegatta.addRaceListener(raceListener, /* Not replicated */ Optional.empty()); } if (!legListeners.containsKey(legIdentifier.getRaceIdentifier())) { TrackedRace trackedRace = racingEventService.getTrackedRace(legIdentifier);