From 0be772ce0e6cf4099da707485819719faa1dcb3f Mon Sep 17 00:00:00 2001 From: MrHua269 Date: Thu, 9 Jul 2026 18:18:59 +0800 Subject: [PATCH] Task bridger --- .../tasks/FlowSchedRunnableTask2CUTask.java | 136 ++++++++++++++ .../tasks/FlowSchedScheduledTask2CUTask.java | 174 ++++++++++++++++++ 2 files changed, 310 insertions(+) create mode 100644 camellia-server/src/main/java/moe/monnairealms/camellia/utils/tasks/FlowSchedRunnableTask2CUTask.java create mode 100644 camellia-server/src/main/java/moe/monnairealms/camellia/utils/tasks/FlowSchedScheduledTask2CUTask.java diff --git a/camellia-server/src/main/java/moe/monnairealms/camellia/utils/tasks/FlowSchedRunnableTask2CUTask.java b/camellia-server/src/main/java/moe/monnairealms/camellia/utils/tasks/FlowSchedRunnableTask2CUTask.java new file mode 100644 index 0000000..7c726e7 --- /dev/null +++ b/camellia-server/src/main/java/moe/monnairealms/camellia/utils/tasks/FlowSchedRunnableTask2CUTask.java @@ -0,0 +1,136 @@ +package moe.monnairealms.camellia.utils.tasks; + +import ca.spottedleaf.concurrentutil.executor.PrioritisedExecutor; +import ca.spottedleaf.concurrentutil.util.Priority; +import com.ishland.c2me.base.common.scheduler.SchedulingManager; +import com.ishland.flowsched.executor.SimpleTask; + +public class FlowSchedRunnableTask2CUTask implements PrioritisedExecutor.PrioritisedTask { + private static final int STATE_CREATED = 0; + private static final int STATE_QUEUED = 1; + private static final int STATE_CANCELLED = 2; + + private volatile int state; + private volatile Priority priority; + + private final SimpleTask task; + private final SchedulingManager worker; + + public FlowSchedRunnableTask2CUTask(Runnable task, SchedulingManager worker, Priority priority) { + this.worker = worker; + this.priority = priority; + this.task = new SimpleTask(task); + } + + public FlowSchedRunnableTask2CUTask(Runnable task, SchedulingManager worker) { + this(task, worker, Priority.NORMAL); + } + + @Override + public PrioritisedExecutor getExecutor() { + throw new UnsupportedOperationException(); + } + + @Override + public boolean queue() { + synchronized (this) { + if (this.state > STATE_CREATED) { + return false; + } + + this.state = STATE_QUEUED; + this.worker.getWorker().schedule(this.task, this.priority.priority * 7); + return true; + } + } + + @Override + public boolean isQueued() { + return this.state == STATE_QUEUED; + } + + @Override + public boolean cancel() { + synchronized (this) { + if (this.state == STATE_CANCELLED || this.state == STATE_QUEUED) { + return false; + } + + this.state = STATE_CANCELLED; + return true; + } + } + + @Override + public Priority getPriority() { + return this.priority; + } + + @Override + public boolean setPriority(Priority priority) { + synchronized (this) { + if (this.state == STATE_CANCELLED || this.state == STATE_QUEUED) { + return false; + } + + this.priority = priority; + this.worker.getWorker().changePriority(this.task, this.priority.priority * 7); + return true; + } + } + + @Override + public boolean setPrioritySubOrderStream(Priority priority, long subOrder, long stream) { + return this.setPriority(priority); + } + + @Override + public boolean raisePriority(Priority priority) { + return this.setPriority(priority); + } + + @Override + public boolean lowerPriority(Priority priority) { + return this.setPriority(priority); + } + + @Override + public boolean execute() { + throw new UnsupportedOperationException(); + } + + @Override + public long getSubOrder() { + throw new UnsupportedOperationException(); + } + + @Override + public boolean setSubOrder(long subOrder) { + throw new UnsupportedOperationException(); + } + + @Override + public boolean raiseSubOrder(long subOrder) { + throw new UnsupportedOperationException(); + } + + @Override + public boolean lowerSubOrder(long subOrder) { + throw new UnsupportedOperationException(); + } + + @Override + public long getStream() { + throw new UnsupportedOperationException(); + } + + @Override + public boolean setStream(long stream) { + throw new UnsupportedOperationException(); + } + + @Override + public PrioritisedExecutor.PriorityState getPriorityState() { + throw new UnsupportedOperationException(); + } +} \ No newline at end of file diff --git a/camellia-server/src/main/java/moe/monnairealms/camellia/utils/tasks/FlowSchedScheduledTask2CUTask.java b/camellia-server/src/main/java/moe/monnairealms/camellia/utils/tasks/FlowSchedScheduledTask2CUTask.java new file mode 100644 index 0000000..2dfddf6 --- /dev/null +++ b/camellia-server/src/main/java/moe/monnairealms/camellia/utils/tasks/FlowSchedScheduledTask2CUTask.java @@ -0,0 +1,174 @@ +package moe.monnairealms.camellia.utils.tasks; + +import ca.spottedleaf.concurrentutil.executor.PrioritisedExecutor; +import ca.spottedleaf.concurrentutil.util.Priority; +import com.ishland.c2me.base.common.scheduler.LockTokenImpl; +import com.ishland.c2me.base.common.scheduler.ScheduledTask; +import com.ishland.c2me.base.common.scheduler.SchedulingManager; +import com.ishland.flowsched.executor.LockToken; +import it.unimi.dsi.fastutil.objects.ObjectArrayList; +import net.minecraft.world.level.ChunkPos; +import org.jetbrains.annotations.Contract; +import org.jetbrains.annotations.NotNull; + +import java.util.concurrent.CompletableFuture; + +public class FlowSchedScheduledTask2CUTask implements PrioritisedExecutor.PrioritisedTask { + private static final int STATE_CREATED = 0; + private static final int STATE_QUEUED = 1; + private static final int STATE_CANCELLED = 2; + + private volatile int state; + private volatile Priority priority; + + private final ScheduledTask task; + private final SchedulingManager worker; + private final ChunkPos pos; + + @Contract("_, _, _, _ -> new") + private static @NotNull ScheduledTask createFlowSchedTaskOf(@NotNull ChunkPos target, int radius, SchedulingManager schedulingManager, Runnable action) { + ObjectArrayList lockTargets = new ObjectArrayList<>((2 * radius + 1) * (2 * radius + 1) + 1); + for (int x = target.x() - radius; x <= target.x() + radius; x++) + for (int z = target.z() - radius; z <= target.z() + radius; z++) + lockTargets.add(new LockTokenImpl(schedulingManager.getId(), ChunkPos.pack(x, z), LockTokenImpl.Usage.WORLDGEN)); + + return new ScheduledTask<>( + target.pack(), + () -> { + action.run(); + return CompletableFuture.completedFuture(null); + }, + lockTargets.toArray(LockToken[]::new)); + } + + public FlowSchedScheduledTask2CUTask(Runnable task, ChunkPos pos, int radius, SchedulingManager worker) { + this(createFlowSchedTaskOf(pos, radius, worker, task), pos, worker, Priority.NORMAL); + } + + public FlowSchedScheduledTask2CUTask(Runnable task, ChunkPos pos, int radius, SchedulingManager worker, Priority priority) { + this(createFlowSchedTaskOf(pos, radius, worker, task), pos, worker, priority); + } + + private FlowSchedScheduledTask2CUTask(ScheduledTask task, ChunkPos pos, SchedulingManager worker, Priority priority) { + this.task = task; + this.pos = pos; + this.worker = worker; + this.priority = priority; + } + + @Override + public PrioritisedExecutor getExecutor() { + throw new UnsupportedOperationException(); + } + + @Override + public boolean queue() { + synchronized (this) { + if (this.state > STATE_CREATED) { + return false; + } + + this.state = STATE_QUEUED; + this.worker.updatePriorityFromLevel(this.pos.pack(), this.priority()); + + this.worker.enqueue(this.task); + return true; + } + } + + @Override + public boolean isQueued() { + return this.state == STATE_QUEUED; + } + + @Override + public boolean cancel() { + synchronized (this) { + if (this.state == STATE_CANCELLED || this.state == STATE_QUEUED) { + return false; + } + + this.state = STATE_CANCELLED; + return true; + } + } + + @Override + public Priority getPriority() { + return this.priority; + } + + private int priority() { + final int invertedPriority = Priority.TOTAL_SCHEDULABLE_PRIORITIES - this.priority.priority - 1; + + return invertedPriority * 7; + } + + @Override + public boolean setPriority(Priority priority) { + synchronized (this) { + if (this.state == STATE_CANCELLED || this.state == STATE_QUEUED) { + return false; + } + + this.priority = priority; + this.worker.updatePriorityFromLevel(this.pos.pack(), this.priority()); + return true; + } + } + + @Override + public boolean setPrioritySubOrderStream(Priority priority, long subOrder, long stream) { + return this.setPriority(priority); + } + + @Override + public boolean raisePriority(Priority priority) { + return this.setPriority(priority); + } + + @Override + public boolean lowerPriority(Priority priority) { + return this.setPriority(priority); + } + + @Override + public boolean execute() { + throw new UnsupportedOperationException(); + } + + @Override + public long getSubOrder() { + throw new UnsupportedOperationException(); + } + + @Override + public boolean setSubOrder(long subOrder) { + throw new UnsupportedOperationException(); + } + + @Override + public boolean raiseSubOrder(long subOrder) { + throw new UnsupportedOperationException(); + } + + @Override + public boolean lowerSubOrder(long subOrder) { + throw new UnsupportedOperationException(); + } + + @Override + public long getStream() { + throw new UnsupportedOperationException(); + } + + @Override + public boolean setStream(long stream) { + throw new UnsupportedOperationException(); + } + + @Override + public PrioritisedExecutor.PriorityState getPriorityState() { + throw new UnsupportedOperationException(); + } +} \ No newline at end of file