Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 28 additions & 24 deletions Essentials/src/main/java/com/earth2me/essentials/AsyncTeleport.java
Original file line number Diff line number Diff line change
Expand Up @@ -139,19 +139,6 @@ public void nowUnsafe(Location loc, TeleportCause cause, CompletableFuture<Boole
paperFuture.exceptionally(future::completeExceptionally);
}

private void runOnMain(final Runnable runnable) throws ExecutionException, InterruptedException {
if (Bukkit.isPrimaryThread()) {
runnable.run();
return;
}
final CompletableFuture<Object> taskLock = new CompletableFuture<>();
Bukkit.getScheduler().runTask(ess, () -> {
runnable.run();
taskLock.complete(new Object());
});
taskLock.get();
}

protected void nowAsync(final IUser teleportee, final ITarget target, final TeleportCause cause, final CompletableFuture<Boolean> future) {
cancel(false);

Expand All @@ -168,14 +155,23 @@ protected void nowAsync(final IUser teleportee, final ITarget target, final Tele
return;
}

try {
runOnMain(() -> teleportee.getBase().eject()); //EntityDismountEvent requires a sync context.
} catch (final ExecutionException | InterruptedException e) {
future.completeExceptionally(e);
return;
// EntityDismountEvent requires the thread that owns the teleportee, which is not necessarily this one
final Runnable ejectAndTeleport = () -> {
teleportee.getBase().eject();
nowAsyncTeleport(teleportee, target, cause, future);
};
if (ess.getTaskScheduler().isOwnedByCurrentThread(teleportee.getBase())) {
ejectAndTeleport.run();
} else {
ess.getTaskScheduler().runEntity(teleportee.getBase(), ejectAndTeleport, () -> future.complete(false), 0);
}
return;
}

nowAsyncTeleport(teleportee, target, cause, future);
}

private void nowAsyncTeleport(final IUser teleportee, final ITarget target, final TeleportCause cause, final CompletableFuture<Boolean> future) {
if (teleportee.isAuthorized("essentials.back.onteleport")) {
teleportee.setLastLocation();
}
Expand All @@ -185,13 +181,12 @@ protected void nowAsync(final IUser teleportee, final ITarget target, final Tele
targetLoc.setX(LocationUtil.getXInsideWorldBorder(targetLoc.getWorld(), targetLoc.getBlockX()));
targetLoc.setZ(LocationUtil.getZInsideWorldBorder(targetLoc.getWorld(), targetLoc.getBlockZ()));
}
PaperLib.getChunkAtAsync(targetLoc.getWorld(), targetLoc.getBlockX() >> 4, targetLoc.getBlockZ() >> 4, true, true).thenAccept(chunk -> {
PaperLib.getChunkAtAsync(targetLoc.getWorld(), targetLoc.getBlockX() >> 4, targetLoc.getBlockZ() >> 4, true, true).thenAccept(chunk -> ess.getTaskScheduler().executeLocation(targetLoc, () -> {
Location loc = targetLoc;
if (LocationUtil.isBlockUnsafeForUser(ess, teleportee, chunk.getWorld(), loc.getBlockX(), loc.getBlockY(), loc.getBlockZ())) {
if (ess.getSettings().isTeleportSafetyEnabled()) {
if (ess.getSettings().isForceDisableTeleportSafety()) {
//The chunk we're teleporting to is 100% going to be loaded here, no need to teleport async.
teleportee.getBase().teleport(loc, cause);
teleportToLoadedChunk(teleportee, loc, cause);
} else {
try {
//There's a chance the safer location is outside the loaded chunk so still teleport async here.
Expand All @@ -207,8 +202,7 @@ protected void nowAsync(final IUser teleportee, final ITarget target, final Tele
}
} else {
if (ess.getSettings().isForceDisableTeleportSafety()) {
//The chunk we're teleporting to is 100% going to be loaded here, no need to teleport async.
teleportee.getBase().teleport(loc, cause);
teleportToLoadedChunk(teleportee, loc, cause);
} else {
if (ess.getSettings().isTeleportToCenterLocation()) {
loc = LocationUtil.getRoundedDestination(loc);
Expand All @@ -218,12 +212,22 @@ protected void nowAsync(final IUser teleportee, final ITarget target, final Tele
}
}
future.complete(true);
}).exceptionally(th -> {
})).exceptionally(th -> {
future.completeExceptionally(th);
return null;
});
}

private void teleportToLoadedChunk(final IUser teleportee, final Location loc, final TeleportCause cause) {
if (ess.getTaskScheduler().isRegionized()) {
// Folia has no synchronous teleports
PaperLib.teleportAsync(teleportee.getBase(), loc, cause);
return;
}
//The chunk we're teleporting to is 100% going to be loaded here, no need to teleport async.
teleportee.getBase().teleport(loc, cause);
}

@Override
public void teleport(final Location loc, final Trade chargeFor, final TeleportCause cause, final CompletableFuture<Boolean> future) {
teleport(teleportOwner, new LocationTarget(loc), chargeFor, cause, future);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@

import net.ess3.api.IEssentials;
import net.ess3.api.IUser;
import net.ess3.provider.TaskSchedulerProvider;
import net.essentialsx.api.v2.events.TeleportWarmupCancelledEvent;
import net.essentialsx.api.v2.events.TeleportWarmupCancelledEvent.CancelReason;

Expand All @@ -31,7 +32,7 @@ public class AsyncTimedTeleport implements Runnable {
private final boolean timer_canMove;
private final Trade timer_chargeFor;
private final TeleportCause timer_cause;
private int timer_task;
private TaskSchedulerProvider.Task timer_task;
private double timer_health;

AsyncTimedTeleport(final IUser user, final IEssentials ess, final AsyncTeleport teleport, final long delay, final IUser teleportUser, final ITarget target, final Trade chargeFor, final TeleportCause cause, final boolean respawn) {
Expand All @@ -55,7 +56,7 @@ public class AsyncTimedTeleport implements Runnable {
this.timer_respawn = respawn;
this.timer_canMove = user.isAuthorized("essentials.teleport.timer.move");

timer_task = ess.runTaskTimerAsynchronously(this, 20, 20).getTaskId();
timer_task = ess.getTaskScheduler().runAsyncTimer(this, 20, 20);

if (future != null) {
this.parentFuture = future;
Expand Down Expand Up @@ -142,16 +143,19 @@ public void run() {
}
}

ess.scheduleSyncDelayedTask(new DelayedTeleportTask());
// If the player has left by now the entity is retired without running the task, so the timer has to stop itself
ess.getTaskScheduler().runEntity(teleportUser.getBase(), new DelayedTeleportTask(), () -> cancelTimer(false), 0);
}

//If we need to cancelTimer a pending teleportPlayer call this method
void cancelTimer(final boolean notifyUser) {
if (timer_task == -1) {
// The timer thread and the thread of the player can both cancel, so read the task once
final TaskSchedulerProvider.Task task = timer_task;
if (task == null) {
return;
}
try {
ess.getServer().getScheduler().cancelTask(timer_task);
task.cancel();

final IUser teleportUser = ess.getUser(this.timer_teleportee);
if (teleportUser != null && teleportUser.getBase() != null) {
Expand All @@ -167,7 +171,7 @@ void cancelTimer(final boolean notifyUser) {
}
}
} finally {
timer_task = -1;
timer_task = null;
}
}
}
19 changes: 10 additions & 9 deletions Essentials/src/main/java/com/earth2me/essentials/Backup.java
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package com.earth2me.essentials;

import net.ess3.api.IEssentials;
import net.ess3.provider.TaskSchedulerProvider;
import org.bukkit.Server;
import org.bukkit.command.CommandSender;

Expand All @@ -18,15 +19,15 @@ public class Backup implements Runnable {
private transient final IEssentials ess;
private final AtomicBoolean pendingShutdown = new AtomicBoolean(false);
private transient boolean running = false;
private transient int taskId = -1;
private transient TaskSchedulerProvider.Task task = null;
private transient boolean active = false;
private transient CompletableFuture<Object> taskLock = null;

public Backup(final IEssentials ess) {
this.ess = ess;
server = ess.getServer();
if (!ess.getOnlinePlayers().isEmpty() || ess.getSettings().isAlwaysRunBackup()) {
ess.runTaskAsynchronously(this::startTask);
ess.getTaskScheduler().runAsync(this::startTask);
}
}

Expand All @@ -36,10 +37,10 @@ public void onPlayerJoin() {

public synchronized void stopTask() {
running = false;
if (taskId != -1) {
server.getScheduler().cancelTask(taskId);
if (task != null) {
task.cancel();
}
taskId = -1;
task = null;
}

private synchronized void startTask() {
Expand All @@ -48,7 +49,7 @@ private synchronized void startTask() {
if (interval < 1200) {
return;
}
taskId = ess.scheduleSyncRepeatingTask(this, interval, interval);
task = ess.getTaskScheduler().runGlobalTimer(this, interval, interval);
running = true;
}
}
Expand Down Expand Up @@ -84,13 +85,13 @@ public void run() {
server.dispatchCommand(cs, "save-all");
server.dispatchCommand(cs, "save-off");

ess.runTaskAsynchronously(() -> {
ess.getTaskScheduler().runAsync(() -> {
try {
final ProcessBuilder childBuilder = new ProcessBuilder(command.split(" "));
childBuilder.redirectErrorStream(true);
childBuilder.directory(ess.getDataFolder().getParentFile().getParentFile());
final Process child = childBuilder.start();
ess.runTaskAsynchronously(() -> {
ess.getTaskScheduler().runAsync(() -> {
try {
try (final BufferedReader reader = new BufferedReader(new InputStreamReader(child.getInputStream()))) {
String line;
Expand Down Expand Up @@ -123,7 +124,7 @@ public void run() {
}

if (!pendingShutdown.get()) {
ess.scheduleSyncDelayedTask(new BackupEnableSaveTask());
ess.getTaskScheduler().runGlobal(new BackupEnableSaveTask());
}
}
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ public CompletableFuture<Void> calculateBalanceTopMapAsync() {
return cacheLock;
}
cacheLock = new CompletableFuture<>();
ess.runTaskAsynchronously(this::calculateBalanceTopMap);
ess.getTaskScheduler().runAsync(this::calculateBalanceTopMap);
return cacheLock;
}

Expand Down
33 changes: 27 additions & 6 deletions Essentials/src/main/java/com/earth2me/essentials/Essentials.java
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@
import com.earth2me.essentials.updatecheck.UpdateChecker;
import com.earth2me.essentials.userstorage.ModernUserMap;
import com.earth2me.essentials.utils.FormatUtil;
import com.earth2me.essentials.utils.ModernPaperEnvironment;
import com.earth2me.essentials.utils.PasteUtil;
import com.earth2me.essentials.utils.VersionUtil;
import io.papermc.lib.PaperLib;
Expand All @@ -70,14 +71,17 @@
import net.ess3.provider.PlayerLocaleProvider;
import net.ess3.provider.ProviderListener;
import net.ess3.provider.ServerStateProvider;
import net.ess3.provider.TaskSchedulerProvider;
import net.ess3.provider.providers.BaseBannerDataProvider;
import net.ess3.provider.providers.BaseInventoryViewProvider;
import net.ess3.provider.providers.BlockMetaSpawnerItemProvider;
import net.ess3.provider.providers.BukkitMaterialTagProvider;
import net.ess3.provider.providers.BukkitSpawnerBlockProvider;
import net.ess3.provider.providers.BukkitTaskSchedulerProvider;
import net.ess3.provider.providers.BukkitTileEntityProvider;
import net.ess3.provider.providers.FixedHeightWorldInfoProvider;
import net.ess3.provider.providers.FlatSpawnEggProvider;
import net.ess3.provider.providers.FoliaTaskSchedulerProvider;
import net.ess3.provider.providers.LegacyBannerDataProvider;
import net.ess3.provider.providers.LegacyBiomeNameProvider;
import net.ess3.provider.providers.LegacyPatternTypeProvider;
Expand Down Expand Up @@ -187,6 +191,7 @@ public class Essentials extends JavaPlugin implements net.ess3.api.IEssentials {
private transient RandomTeleport randomTeleport;
private transient UpdateChecker updateChecker;
private transient AdventureFacet adventureFacet;
private transient TaskSchedulerProvider taskScheduler;

static {
EconomyLayers.init();
Expand Down Expand Up @@ -261,6 +266,10 @@ public void onEnable() {
getLogger().severe(getAdventureFacet().miniToLegacy(tlLiteral("serverSnapshot")));
}

if (PaperLib.isPaper() && !PaperLib.isVersion(13) && VersionUtil.getServerBukkitVersion().isHigherThanOrEqualTo(VersionUtil.v1_13_0_R01)) {
PaperLib.setCustomEnvironment(new ModernPaperEnvironment());
}

final PluginManager pm = getServer().getPluginManager();
for (final Plugin plugin : pm.getPlugins()) {
if (plugin.getDescription().getName().startsWith("Essentials") && !plugin.getDescription().getVersion().equals(this.getDescription().getVersion()) && !plugin.getDescription().getName().equals("EssentialsAntiCheat")) {
Expand Down Expand Up @@ -327,9 +336,6 @@ public void onEnable() {
confList.add(jails);
execTimer.mark("Init(Jails)");

EconomyLayers.onEnable(this);
execTimer.mark("Init(EconomyLayers)");

// Spawner item provider only uses one, but it's here for legacy...
providerFactory.registerProvider(BlockMetaSpawnerItemProvider.class);

Expand Down Expand Up @@ -405,9 +411,16 @@ public void onEnable() {
// Tick Count Provider
providerFactory.registerProvider(PaperTickCountProvider.class);

// Task Scheduler Provider
providerFactory.registerProvider(BukkitTaskSchedulerProvider.class, FoliaTaskSchedulerProvider.class);

if (!TESTING) {
providerFactory.finalizeRegistration();
}
taskScheduler = TESTING ? new BukkitTaskSchedulerProvider(this) : provider(TaskSchedulerProvider.class);

EconomyLayers.onEnable(this);
execTimer.mark("Init(EconomyLayers)");

// Event Providers
if (PaperLib.isPaper()) {
Expand Down Expand Up @@ -436,7 +449,7 @@ public void onEnable() {
alternativeCommandsHandler = new AlternativeCommandsHandler(this);

timer = new EssentialsTimer(this);
scheduleSyncRepeatingTask(timer, 1000, 50);
taskScheduler.runGlobalTimer(timer, 1000, 50);

Economy.setEss(this);
execTimer.mark("RegHandler");
Expand All @@ -447,7 +460,7 @@ public void onEnable() {

if (!TESTING) {
updateChecker = new UpdateChecker(this);
runTaskAsynchronously(() -> {
taskScheduler.runAsync(() -> {
getLogger().log(Level.INFO, getAdventureFacet().miniToLegacy(tlLiteral("versionFetching")));
for (final ComponentHolder component : updateChecker.getVersionMessages(false, true, new CommandSource(this, Bukkit.getConsoleSender()))) {
getLogger().log(getSettings().isUpdateCheckEnabled() ? Level.WARNING : Level.INFO, getAdventureFacet().adventureToLegacy(component));
Expand Down Expand Up @@ -589,7 +602,9 @@ public void onDisable() {

EssentialsConfiguration.shutdownExecutor();
PasteUtil.shutdownExecutor();
getServer().getScheduler().cancelTasks(this);
if (taskScheduler != null) {
taskScheduler.cancelAll();
}

HandlerList.unregisterAll(this);
}
Expand Down Expand Up @@ -911,10 +926,16 @@ public void showError(final CommandSource sender, final Throwable exception, fin
}

@Override
@Deprecated
public BukkitScheduler getScheduler() {
return this.getServer().getScheduler();
}

@Override
public TaskSchedulerProvider getTaskScheduler() {
return taskScheduler;
}

@Override
public List<Player> getJailedPlayers() {
return getUsers().getAllUserUUIDs().stream().map(this::getUser).filter(Objects::nonNull).filter(User::isJailed).map(User::getBase).filter(Objects::nonNull).collect(Collectors.toList());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ public void onBlockPlace(final BlockPlaceEvent event) {

final User user = ess.getUser(event.getPlayer());
if (user.hasUnlimited(is) && user.getBase().getGameMode() == GameMode.SURVIVAL) {
ess.scheduleSyncDelayedTask(() -> {
ess.getTaskScheduler().runEntity(user.getBase(), () -> {
if (is != null && is.getType() != null && !MaterialUtil.isAir(is.getType())) {
final ItemStack cloneIs = is.clone();
cloneIs.setAmount(1);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,7 @@ public void run() {
}
}

ess.scheduleSyncDelayedTask(new PowerToolInteractTask());
ess.getTaskScheduler().runEntity(attacker.getBase(), new PowerToolInteractTask());

event.setCancelled(true);
return;
Expand Down
Loading