This commit is contained in:
Len
2022-10-19 17:31:20 +02:00
parent dcdba51f14
commit 2ef0fbf727
8 changed files with 16 additions and 0 deletions
@@ -0,0 +1,76 @@
package com.alttd;
import com.google.common.io.ByteArrayDataOutput;
import com.google.common.io.ByteStreams;
import org.bukkit.scheduler.BukkitRunnable;
import java.util.HashSet;
import java.util.UUID;
public class DataLock {
private static DataLock instance = null;
public static DataLock getInstance() {
if (instance == null)
instance = new DataLock();
return instance;
}
private final PluginMessageListener pluginMessageListener;
private final DataLockLib plugin;
private final Idempotency activeRequests;
private DataLock() {
pluginMessageListener = new PluginMessageListener();
plugin = DataLockLib.getInstance();
activeRequests = new Idempotency();
new RepeatRequest().runTaskTimerAsynchronously(plugin, 15 * 20, 15 * 20);
//Run repeat request task every 15 seconds
}
private void sendPluginMessage(RequestType requestType, IdempotencyData idempotencyData) {
ByteArrayDataOutput out = ByteStreams.newDataOutput();
out.writeUTF(requestType.subChannel);
out.writeUTF(idempotencyData.data());
plugin.getServer().sendPluginMessage(plugin, idempotencyData.channel(), out.toByteArray());
}
private final HashSet<String> activeChannels = new HashSet<>();
protected boolean removeActiveRequest(RequestType requestType, IdempotencyData idempotencyData) {
return activeRequests.removeIdempotencyData(requestType, idempotencyData);
}
public synchronized void registerChannel(String channel) {
activeChannels.add(channel);
plugin.getServer().getMessenger().registerOutgoingPluginChannel(plugin, channel);
plugin.getServer().getMessenger().registerIncomingPluginChannel(plugin, channel, pluginMessageListener);
}
public synchronized boolean isActiveChannel(String channel) {
return activeChannels.contains(channel);
}
public void tryLock(String channel, String data) {
IdempotencyData idempotencyData = new IdempotencyData(channel, data, UUID.randomUUID());
activeRequests.putIdempotencyData(RequestType.TRY_LOCK, idempotencyData);
sendPluginMessage(RequestType.TRY_LOCK, idempotencyData);
}
public void tryUnlock(String channel, String data) {
IdempotencyData idempotencyData = new IdempotencyData(channel, data, UUID.randomUUID());
activeRequests.putIdempotencyData(RequestType.TRY_UNLOCK, idempotencyData);
sendPluginMessage(RequestType.TRY_UNLOCK, idempotencyData);
}
private class RepeatRequest extends BukkitRunnable {
@Override
public void run() {
for (RequestType requestType : RequestType.values()) {
for (IdempotencyData next : activeRequests.getIdempotencyData(requestType)) {
sendPluginMessage(requestType, next);
}
}
}
}
}
@@ -0,0 +1,27 @@
package com.alttd;
import org.bukkit.plugin.java.JavaPlugin;
public class DataLockLib extends JavaPlugin {
private static DataLockLib instance;
protected static DataLockLib getInstance() {
return instance;
}
@Override
public void onLoad() {
instance = this;
}
@Override
public void onEnable() {
}
@Override
public void onDisable() {
}
}
@@ -0,0 +1,51 @@
package com.alttd;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Set;
class Idempotency {
private final HashMap<RequestType, HashSet<IdempotencyData>> idempotencyMap;
protected Idempotency() {
idempotencyMap = new HashMap<>();
}
private HashSet<IdempotencyData> getIdempotencySet(RequestType requestType) {
return idempotencyMap.getOrDefault(requestType, new HashSet<>());
}
private void putIdempotencySet(RequestType requestType, HashSet<IdempotencyData> idempotencySet) {
idempotencyMap.put(requestType, idempotencySet);
}
/**
* Add IdempotencyData to the list of active queries
* @param idempotencyData Data to add to the set
* @return true if entry did not exist yet
*/
protected synchronized boolean putIdempotencyData(RequestType requestType, IdempotencyData idempotencyData) {
HashSet<IdempotencyData> idempotencySet = getIdempotencySet(requestType);
boolean result = idempotencySet.add(idempotencyData);
putIdempotencySet(requestType, idempotencySet);
return result;
}
/**
* Remove IdempotencyData from the list of active queries
* @param idempotencyData Data to remove from the set
* @return True if the data that was requested to be removed was in the set and was removed
*/
protected synchronized boolean removeIdempotencyData(RequestType requestType, IdempotencyData idempotencyData) {
HashSet<IdempotencyData> idempotencySet = getIdempotencySet(requestType);
boolean result = idempotencySet.remove(idempotencyData);
putIdempotencySet(requestType, idempotencySet);
return result;
}
protected synchronized Set<IdempotencyData> getIdempotencyData(RequestType requestType) {
return Collections.unmodifiableSet(getIdempotencySet(requestType));
}
}
@@ -0,0 +1,25 @@
package com.alttd;
import java.util.UUID;
record IdempotencyData(String channel, String data, UUID idempotencyToken) {
@Override
public String toString() {
return "Channel: [" + channel + "] Data: [" + data + "] Idempotency Token: [" + idempotencyToken + "]";
}
@Override
public String channel() {
return channel;
}
@Override
public String data() {
return data;
}
@Override
public UUID idempotencyToken() {
return idempotencyToken;
}
}
@@ -0,0 +1,42 @@
package com.alttd;
import org.bukkit.event.Event;
import org.bukkit.event.HandlerList;
import org.jetbrains.annotations.NotNull;
public class LockResponseEvent extends Event {
private final HandlerList handlers = new HandlerList();
private final String channel;
private final ResponseType responseType;
private final String data;
private final boolean result;
protected LockResponseEvent(boolean isAsync, String channel, ResponseType responseType, String data, boolean result) {
super(isAsync);
this.channel = channel;
this.responseType = responseType;
this.data = data;
this.result = result;
}
public String getChannel() {
return channel;
}
public ResponseType getResponseType() {
return responseType;
}
public String getData() {
return data;
}
public boolean getResult() {
return result;
}
public @NotNull HandlerList getHandlers() {
return handlers;
}
}
@@ -0,0 +1,68 @@
package com.alttd;
import com.google.common.io.ByteArrayDataInput;
import com.google.common.io.ByteStreams;
import org.bukkit.entity.Player;
import org.bukkit.scheduler.BukkitRunnable;
import org.jetbrains.annotations.NotNull;
import java.util.UUID;
class PluginMessageListener implements org.bukkit.plugin.messaging.PluginMessageListener {
private final DataLock dataLock;
private final Idempotency alreadyReceived;
PluginMessageListener() {
this.dataLock = DataLock.getInstance();
this.alreadyReceived = new Idempotency();
}
@Override
public void onPluginMessageReceived(@NotNull String channel, @NotNull Player player, byte[] bytes) {
if (!dataLock.isActiveChannel(channel)) {
return;
}
ByteArrayDataInput in = ByteStreams.newDataInput(bytes);
String data = in.readUTF();
boolean result = in.readBoolean();
UUID idempotency = UUID.fromString(in.readUTF());
IdempotencyData idempotencyData = new IdempotencyData(channel, data, idempotency);
new BukkitRunnable() {
@Override
public void run() {
switch (in.readUTF()) {
case "try-lock-result" -> {
if (!alreadyReceived.putIdempotencyData(RequestType.TRY_LOCK, idempotencyData))
return;
DataLock.getInstance().removeActiveRequest(RequestType.TRY_LOCK, idempotencyData);
new LockResponseEvent(true, channel, ResponseType.TRY_LOCK_RESULT, data, result);
}
case "queue-lock-failed" -> {
if (!alreadyReceived.putIdempotencyData(RequestType.TRY_LOCK, idempotencyData))
return;
DataLock.getInstance().removeActiveRequest(RequestType.TRY_LOCK, idempotencyData);
new LockResponseEvent(true, channel, ResponseType.QUEUE_LOCK_FAILED, data, result);
}
case "try-unlock-result" -> {
if (!alreadyReceived.putIdempotencyData(RequestType.TRY_UNLOCK, idempotencyData))
return;
DataLock.getInstance().removeActiveRequest(RequestType.TRY_UNLOCK, idempotencyData);
new LockResponseEvent(true, channel, ResponseType.TRY_UNLOCK_RESULT, data, result);
}
case "locked-queue-lock" -> {
if (!alreadyReceived.putIdempotencyData(RequestType.TRY_LOCK, idempotencyData))
return;
DataLock.getInstance().removeActiveRequest(RequestType.TRY_LOCK, idempotencyData);
new LockResponseEvent(true, channel, ResponseType.LOCKED_QUEUE_LOCK, data, result);
}
case "check-lock-result" -> {
if (!alreadyReceived.putIdempotencyData(RequestType.CHECK_LOCK, idempotencyData))
return;
DataLock.getInstance().removeActiveRequest(RequestType.CHECK_LOCK, idempotencyData);
new LockResponseEvent(true, channel, ResponseType.CHECK_LOCK_RESULT, data, result);
}
}
}
}.runTaskAsynchronously(DataLockLib.getInstance());
}
}
@@ -0,0 +1,13 @@
package com.alttd;
enum RequestType {
TRY_LOCK("try-lock"),
TRY_UNLOCK("try-unlock"),
CHECK_LOCK("check-lock");
String subChannel;
RequestType(String subChannel) {
this.subChannel = subChannel;
}
}
@@ -0,0 +1,9 @@
package com.alttd;
public enum ResponseType {
TRY_LOCK_RESULT,
QUEUE_LOCK_FAILED,
TRY_UNLOCK_RESULT,
LOCKED_QUEUE_LOCK,
CHECK_LOCK_RESULT
}