Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<VoteTimeQueue> 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.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<VoteTimeQueue> timeChangeQueue = new ConcurrentLinkedQueue<>();

private VotingPluginMain plugin;
private final AtomicBoolean retryPending = new AtomicBoolean();
private final AtomicInteger retryAttempts = new AtomicInteger();
private final Set<VoteTimeQueue> completedAwaitingPersistence =
Collections.newSetFromMap(new IdentityHashMap<>());

/**
* Constructs a new TimeQueueHandler.
Expand All @@ -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);
}

Expand All @@ -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();
Expand Down Expand Up @@ -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<VoteTimeQueue> pending = new ArrayList<>(timeChangeQueue);
pending.removeAll(completedAwaitingPersistence);
plugin.getServerData().replaceTimedVoteCache(pending);
}

private boolean persistWithout(VoteTimeQueue completed) {
ArrayList<VoteTimeQueue> 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;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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;
Expand Down Expand Up @@ -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<java.util.Collection<com.bencodez.votingplugin.timequeue.VoteTimeQueue>> 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);
Expand Down
Loading