diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/data/ServerData.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/data/ServerData.java index a255b8d58..60776e1f3 100644 --- a/VotingPlugin/src/main/java/com/bencodez/votingplugin/data/ServerData.java +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/data/ServerData.java @@ -4,6 +4,7 @@ import java.util.ArrayList; import java.util.HashSet; import java.util.List; +import java.util.Collection; import java.util.Locale; import java.util.Set; import java.util.UUID; @@ -110,6 +111,24 @@ public void clearTimedVoteCache() { saveData(); } + /** + * Replaces the timed vote cache and persists the complete snapshot with one save. + * + * @param votes pending timed votes in replay order + */ + public synchronized void replaceTimedVoteCache(Collection votes) { + ConfigurationSection data = getData(); + data.set("TimedVoteCache", null); + int index = 0; + for (VoteTimeQueue vote : votes) { + String path = "TimedVoteCache." + index++; + data.set(path + ".Name", vote.getName()); + data.set(path + ".Service", vote.getService()); + data.set(path + ".Time", vote.getTime()); + } + saveData(); + } + /** * Gets the auto cached placeholders. * diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/timequeue/TimeQueueHandler.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/timequeue/TimeQueueHandler.java index 3bf50342a..4926029d1 100644 --- a/VotingPlugin/src/main/java/com/bencodez/votingplugin/timequeue/TimeQueueHandler.java +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/timequeue/TimeQueueHandler.java @@ -3,6 +3,10 @@ import java.time.LocalDateTime; import java.time.ZoneId; import java.util.Queue; +import java.util.ArrayList; +import java.util.Collections; +import java.util.IdentityHashMap; +import java.util.Set; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -25,12 +29,15 @@ */ public class TimeQueueHandler implements Listener { private static final long TICKS_PER_SECOND = 20L; + static final int MAX_QUEUED_VOTES = 4096; @Getter private Queue timeChangeQueue = new ConcurrentLinkedQueue<>(); private VotingPluginMain plugin; private final AtomicBoolean retryPending = new AtomicBoolean(); private final AtomicInteger retryAttempts = new AtomicInteger(); + private final Set completedAwaitingPersistence = + Collections.newSetFromMap(new IdentityHashMap<>()); /** * Constructs a new TimeQueueHandler. @@ -48,20 +55,33 @@ public TimeQueueHandler(VotingPluginMain plugin) { * @param voteUsername the voter username * @param voteSiteName the vote site name */ - public void addVote(String voteUsername, String voteSiteName) { + public synchronized void addVote(String voteUsername, String voteSiteName) { + if (timeChangeQueue.size() >= MAX_QUEUED_VOTES) { + plugin.getLogger().severe("Time-change vote queue is full; rejecting vote instead of expanding durable storage"); + return; + } timeChangeQueue.add(new VoteTimeQueue(voteUsername, voteSiteName, LocalDateTime.now().atZone(ZoneId.systemDefault()).toInstant().toEpochMilli())); + persistQueueSnapshot(); } /** * Loads cached votes from server data and schedules queue processing. */ public void load() { + boolean truncated = false; for (String str : plugin.getServerData().getTimedVoteCacheKeys()) { + if (timeChangeQueue.size() >= MAX_QUEUED_VOTES) { + truncated = true; + break; + } ConfigurationSection data = plugin.getServerData().getTimedVoteCacheSection(str); timeChangeQueue .add(new VoteTimeQueue(data.getString("Name"), data.getString("Service"), data.getLong("Time"))); } + if (truncated) { + plugin.getLogger().severe("Timed vote recovery exceeded the bounded queue; excess persisted votes were not loaded"); + } scheduleQueueProcessing(120, TimeUnit.SECONDS); } @@ -76,12 +96,7 @@ public void postTimeChange(DateChangedEvent event) { } private void scheduleQueueProcessing(long delay, TimeUnit unit) { - boolean admitted = VoteTaskAdmission.trySchedule(plugin.getVoteTimer(), () -> { - // Clear only after the bounded executor has admitted the task. If it is - // rejected, shutdown persistence can still recover the in-memory queue. - plugin.getServerData().clearTimedVoteCache(); - processQueue(); - }, delay, unit); + boolean admitted = VoteTaskAdmission.trySchedule(plugin.getVoteTimer(), this::processQueue, delay, unit); if (!admitted) { plugin.getLogger().warning("Unable to schedule time-queue processing because vote processing is busy; queued votes were retained."); scheduleRetry(); @@ -109,33 +124,61 @@ private void scheduleRetry() { /** * Processes all votes in the queue. */ - public void processQueue() { - while (getTimeChangeQueue().size() > 0) { - VoteTimeQueue vote = getTimeChangeQueue().remove(); - PlayerVoteEvent voteEvent = new PlayerVoteEvent( - plugin.getVoteSiteManager().getVoteSite(plugin.getVoteSiteManager().getVoteSiteName(true, vote.getService()), true), vote.getName(), - vote.getService(), true); - voteEvent.setTime(vote.getTime()); - plugin.getServer().getPluginManager().callEvent(voteEvent); - - if (voteEvent.isCancelled()) { - plugin.debug("Vote cancelled"); + public synchronized void processQueue() { + while (true) { + VoteTimeQueue vote = getTimeChangeQueue().peek(); + if (vote == null) return; + if (!completedAwaitingPersistence.contains(vote)) { + PlayerVoteEvent voteEvent = new PlayerVoteEvent( + plugin.getVoteSiteManager().getVoteSite(plugin.getVoteSiteManager().getVoteSiteName(true, vote.getService()), true), vote.getName(), + vote.getService(), true); + voteEvent.setTime(vote.getTime()); + try { + plugin.getServer().getPluginManager().callEvent(voteEvent); + } catch (RuntimeException failure) { + plugin.getLogger().warning("Unable to process queued time-change vote; retaining it for retry"); + plugin.debug(failure); + scheduleRetry(); + return; + } + completedAwaitingPersistence.add(vote); + if (voteEvent.isCancelled()) plugin.debug("Vote cancelled"); + } + + if (!persistWithout(vote)) { + scheduleRetry(); return; } + getTimeChangeQueue().remove(vote); + completedAwaitingPersistence.remove(vote); } } /** * Saves pending votes to server data. */ - public void save() { - if (!timeChangeQueue.isEmpty()) { - int num = 0; - for (VoteTimeQueue vote : timeChangeQueue) { - plugin.getServerData().addTimeVoted(num, vote); - num++; - } - } + public synchronized void save() { + persistQueueSnapshot(); timeChangeQueue.clear(); + completedAwaitingPersistence.clear(); + } + + private void persistQueueSnapshot() { + ArrayList pending = new ArrayList<>(timeChangeQueue); + pending.removeAll(completedAwaitingPersistence); + plugin.getServerData().replaceTimedVoteCache(pending); + } + + private boolean persistWithout(VoteTimeQueue completed) { + ArrayList remaining = new ArrayList<>(timeChangeQueue); + remaining.remove(completed); + try { + plugin.getServerData().replaceTimedVoteCache(remaining); + return true; + } catch (RuntimeException failure) { + plugin.getLogger().warning("Unable to persist completed time-change vote retirement; retrying storage only"); + plugin.debug(failure); + return false; + } } } diff --git a/VotingPlugin/src/test/java/com/bencodez/votingplugin/tests/timequeue/TimeQueueHandlerRejectionTest.java b/VotingPlugin/src/test/java/com/bencodez/votingplugin/tests/timequeue/TimeQueueHandlerRejectionTest.java index 37ba7a49a..b1d30c361 100644 --- a/VotingPlugin/src/test/java/com/bencodez/votingplugin/tests/timequeue/TimeQueueHandlerRejectionTest.java +++ b/VotingPlugin/src/test/java/com/bencodez/votingplugin/tests/timequeue/TimeQueueHandlerRejectionTest.java @@ -6,10 +6,12 @@ import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.RETURNS_DEEP_STUBS; +import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.when; import static org.mockito.Mockito.reset; @@ -20,13 +22,17 @@ import java.util.logging.Logger; import org.bukkit.configuration.ConfigurationSection; +import org.bukkit.Server; +import org.bukkit.plugin.PluginManager; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import com.bencodez.advancedcore.api.time.events.DateChangedEvent; import com.bencodez.votingplugin.VotingPluginMain; import com.bencodez.votingplugin.data.ServerData; +import com.bencodez.votingplugin.events.PlayerVoteEvent; import com.bencodez.votingplugin.timequeue.TimeQueueHandler; +import com.bencodez.votingplugin.timequeue.VoteTimeQueue; class TimeQueueHandlerRejectionTest { private VotingPluginMain plugin; @@ -73,6 +79,93 @@ void dateChangeRejectionDoesNotEscapeBukkitEventHandler() { verify(logger, org.mockito.Mockito.atLeastOnce()).warning(anyString()); } + @Test + void fullTimeQueueRejectsInsteadOfExpandingDurableStorage() { + TimeQueueHandler handler = new TimeQueueHandler(plugin); + handler.getTimeChangeQueue().clear(); + for (int i = 0; i < 4096; i++) { + handler.getTimeChangeQueue().add(new VoteTimeQueue("Player" + i, "example.org", i + 1L)); + } + org.mockito.Mockito.clearInvocations(serverData, logger); + + handler.addVote("Overflow", "example.org"); + + assertEquals(4096, handler.getTimeChangeQueue().size()); + verify(serverData, never()).replaceTimedVoteCache(any()); + verify(logger).severe(org.mockito.ArgumentMatchers.contains("queue is full")); + } + + @Test + void newlyQueuedVoteIsPersistedImmediately() { + TimeQueueHandler handler = new TimeQueueHandler(plugin); + org.mockito.Mockito.clearInvocations(serverData); + + handler.addVote("Alex", "example.org"); + + org.mockito.ArgumentCaptor> snapshot = + org.mockito.ArgumentCaptor.forClass(java.util.Collection.class); + verify(serverData).replaceTimedVoteCache(snapshot.capture()); + assertEquals(2, snapshot.getValue().size()); + } + + @Test + void processingFailureRetainsDurableQueueHead() { + TimeQueueHandler handler = new TimeQueueHandler(plugin); + Server server = mock(Server.class); + PluginManager manager = mock(PluginManager.class); + when(plugin.getServer()).thenReturn(server); + when(server.getPluginManager()).thenReturn(manager); + when(plugin.getVoteSiteManager().getVoteSiteName(true, "example.org")).thenReturn("example.org"); + doThrow(new IllegalStateException("listener failed")).when(manager).callEvent(any(PlayerVoteEvent.class)); + org.mockito.Mockito.clearInvocations(serverData); + + handler.processQueue(); + + assertEquals(1, handler.getTimeChangeQueue().size()); + verify(serverData, never()).replaceTimedVoteCache(any()); + } + + @Test + void completedVoteDoesNotRepeatEffectsWhenRetirementPersistenceRetries() { + TimeQueueHandler handler = new TimeQueueHandler(plugin); + Server server = mock(Server.class); + PluginManager manager = mock(PluginManager.class); + when(plugin.getServer()).thenReturn(server); + when(server.getPluginManager()).thenReturn(manager); + when(plugin.getVoteSiteManager().getVoteSiteName(true, "example.org")).thenReturn("example.org"); + doThrow(new IllegalStateException("save failed")).doNothing() + .when(serverData).replaceTimedVoteCache(any()); + + handler.processQueue(); + assertEquals(1, handler.getTimeChangeQueue().size()); + verify(manager, times(1)).callEvent(any(PlayerVoteEvent.class)); + + handler.processQueue(); + assertEquals(0, handler.getTimeChangeQueue().size()); + verify(manager, times(1)).callEvent(any(PlayerVoteEvent.class)); + } + + @Test + void cancelledVoteDoesNotStrandLaterQueuedVotes() { + TimeQueueHandler handler = new TimeQueueHandler(plugin); + handler.addVote("Alex", "second.example.org"); + Server server = mock(Server.class); + PluginManager manager = mock(PluginManager.class); + when(plugin.getServer()).thenReturn(server); + when(server.getPluginManager()).thenReturn(manager); + when(plugin.getVoteSiteManager().getVoteSiteName(true, anyString())).thenAnswer(invocation -> invocation.getArgument(1)); + doAnswer(invocation -> { + ((PlayerVoteEvent) invocation.getArgument(0)).setCancelled(true); + return null; + }).when(manager).callEvent(any(PlayerVoteEvent.class)); + org.mockito.Mockito.clearInvocations(serverData); + + handler.processQueue(); + + assertEquals(0, handler.getTimeChangeQueue().size()); + verify(manager, times(2)).callEvent(any(PlayerVoteEvent.class)); + } + @Test void rejectedProcessingSchedulesOneBoundedRetry() { TimeQueueHandler handler = new TimeQueueHandler(plugin);