-
Notifications
You must be signed in to change notification settings - Fork 47
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Ensure all locks are available before real acquisition and task execu…
…tion Closes #2
- Loading branch information
Showing
4 changed files
with
121 additions
and
5 deletions.
There are no files selected for viewing
15 changes: 14 additions & 1 deletion
15
src/main/java/org/yatopiamc/c2me/common/threading/GlobalExecutors.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,14 +1,27 @@ | ||
package org.yatopiamc.c2me.common.threading; | ||
|
||
import com.google.common.util.concurrent.ThreadFactoryBuilder; | ||
import org.jetbrains.annotations.NotNull; | ||
|
||
import java.util.concurrent.ScheduledThreadPoolExecutor; | ||
import java.util.concurrent.ThreadFactory; | ||
import java.util.concurrent.atomic.AtomicReference; | ||
|
||
public class GlobalExecutors { | ||
|
||
public static final ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor( | ||
1, | ||
new ThreadFactoryBuilder().setNameFormat("C2ME scheduler").setDaemon(true).setPriority(Thread.NORM_PRIORITY - 1).build() | ||
new ThreadFactoryBuilder().setNameFormat("C2ME scheduler").setDaemon(true).setPriority(Thread.NORM_PRIORITY - 1).setThreadFactory(r -> { | ||
final Thread thread = new Thread(r); | ||
GlobalExecutors.schedulerThread.set(thread); | ||
return thread; | ||
}).build() | ||
); | ||
private static final AtomicReference<Thread> schedulerThread = new AtomicReference<>(); | ||
|
||
public static void ensureSchedulerThread() { | ||
if (Thread.currentThread() != schedulerThread.get()) | ||
throw new IllegalStateException("Not on scheduler thread"); | ||
} | ||
|
||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
70 changes: 70 additions & 0 deletions
70
src/main/java/org/yatopiamc/c2me/common/util/AsyncCombinedLock.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,70 @@ | ||
package org.yatopiamc.c2me.common.util; | ||
|
||
import com.ibm.asyncutil.locks.AsyncLock; | ||
import org.yatopiamc.c2me.common.threading.GlobalExecutors; | ||
|
||
import java.util.Optional; | ||
import java.util.Set; | ||
import java.util.concurrent.CompletableFuture; | ||
import java.util.function.Function; | ||
import java.util.stream.Collectors; | ||
|
||
public class AsyncCombinedLock { | ||
|
||
private final Set<AsyncLock> lockHandles; | ||
private final CompletableFuture<AsyncLock.LockToken> future = new CompletableFuture<>(); | ||
|
||
public AsyncCombinedLock(Set<AsyncLock> lockHandles) { | ||
this.lockHandles = Set.copyOf(lockHandles); | ||
GlobalExecutors.scheduler.execute(this::tryAcquire); | ||
} | ||
|
||
private void tryAcquire() { | ||
GlobalExecutors.ensureSchedulerThread(); | ||
final Set<LockEntry> tryLocks = lockHandles.stream().map(lock -> new LockEntry(lock, lock.tryLock())).collect(Collectors.toUnmodifiableSet()); | ||
if (tryLocks.stream().allMatch(lockEntry -> lockEntry.lockToken.isPresent())) { | ||
future.complete(new CombinedLockToken(tryLocks.stream().flatMap(lockEntry -> lockEntry.lockToken.stream()).collect(Collectors.toUnmodifiableSet()))); | ||
} else { | ||
tryLocks.stream().flatMap(lockEntry -> lockEntry.lockToken.stream()).forEach(AsyncLock.LockToken::releaseLock); | ||
tryLocks.stream().unordered().filter(lockEntry -> lockEntry.lockToken.isEmpty()).findFirst().ifPresentOrElse(lockEntry -> | ||
lockEntry.lock.acquireLock().thenCompose(lockToken -> { | ||
lockToken.releaseLock(); | ||
return CompletableFuture.runAsync(this::tryAcquire, GlobalExecutors.scheduler); | ||
}), this::tryAcquire); | ||
} | ||
} | ||
|
||
public CompletableFuture<AsyncLock.LockToken> getFuture() { | ||
return future.thenApply(Function.identity()); | ||
} | ||
|
||
@SuppressWarnings("OptionalUsedAsFieldOrParameterType") | ||
private static class LockEntry { | ||
public final AsyncLock lock; | ||
public final Optional<AsyncLock.LockToken> lockToken; | ||
|
||
private LockEntry(AsyncLock lock, Optional<AsyncLock.LockToken> lockToken) { | ||
this.lock = lock; | ||
this.lockToken = lockToken; | ||
} | ||
} | ||
|
||
private static class CombinedLockToken implements AsyncLock.LockToken { | ||
|
||
private final Set<AsyncLock.LockToken> delegates; | ||
|
||
private CombinedLockToken(Set<AsyncLock.LockToken> delegates) { | ||
this.delegates = Set.copyOf(delegates); | ||
} | ||
|
||
@Override | ||
public void releaseLock() { | ||
delegates.forEach(AsyncLock.LockToken::releaseLock); | ||
} | ||
|
||
@Override | ||
public void close() { | ||
this.releaseLock(); | ||
} | ||
} | ||
} |
29 changes: 29 additions & 0 deletions
29
src/main/java/org/yatopiamc/c2me/common/util/AsyncNamedLockDelegateAsyncLock.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,29 @@ | ||
package org.yatopiamc.c2me.common.util; | ||
|
||
import com.ibm.asyncutil.locks.AsyncLock; | ||
import com.ibm.asyncutil.locks.AsyncNamedLock; | ||
|
||
import java.util.Objects; | ||
import java.util.Optional; | ||
import java.util.concurrent.CompletionStage; | ||
|
||
public class AsyncNamedLockDelegateAsyncLock<T> implements AsyncLock { | ||
|
||
private final AsyncNamedLock<T> delegate; | ||
private final T name; | ||
|
||
public AsyncNamedLockDelegateAsyncLock(AsyncNamedLock<T> delegate, T name) { | ||
this.delegate = Objects.requireNonNull(delegate); | ||
this.name = name; | ||
} | ||
|
||
@Override | ||
public CompletionStage<LockToken> acquireLock() { | ||
return delegate.acquireLock(name); | ||
} | ||
|
||
@Override | ||
public Optional<LockToken> tryLock() { | ||
return delegate.tryLock(name); | ||
} | ||
} |