Use spring-cloud-build-tools for checkstyle validation
This commit is contained in:
6
pom.xml
6
pom.xml
@@ -150,8 +150,8 @@
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>org.springframework.cloud</groupId>
|
||||
<artifactId>spring-cloud-stream-tools</artifactId>
|
||||
<version>1.1.0.BUILD-SNAPSHOT</version>
|
||||
<artifactId>spring-cloud-build-tools</artifactId>
|
||||
<version>1.2.0.RELEASE</version>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
<executions>
|
||||
@@ -160,8 +160,6 @@
|
||||
<phase>validate</phase>
|
||||
<configuration>
|
||||
<configLocation>checkstyle.xml</configLocation>
|
||||
<headerLocation>checkstyle-header.txt</headerLocation>
|
||||
<suppressionsLocation>checkstyle-suppressions.xml</suppressionsLocation>
|
||||
<encoding>UTF-8</encoding>
|
||||
<consoleOutput>true</consoleOutput>
|
||||
<failsOnError>true</failsOnError>
|
||||
|
||||
@@ -50,297 +50,300 @@ import reactor.util.concurrent.OpenHashSet;
|
||||
* @author Stephane Maldini
|
||||
*/
|
||||
final class ElasticScheduler implements Scheduler {
|
||||
static final AtomicLong COUNTER = new AtomicLong();
|
||||
static final AtomicLong COUNTER = new AtomicLong();
|
||||
|
||||
static final ThreadFactory EVICTOR_FACTORY = r -> {
|
||||
Thread t = new Thread(r, "elastic-evictor-" + COUNTER.incrementAndGet());
|
||||
t.setDaemon(true);
|
||||
return t;
|
||||
};
|
||||
static final ThreadFactory EVICTOR_FACTORY = r -> {
|
||||
Thread t = new Thread(r, "elastic-evictor-" + COUNTER.incrementAndGet());
|
||||
t.setDaemon(true);
|
||||
return t;
|
||||
};
|
||||
|
||||
final ThreadFactory factory;
|
||||
|
||||
final int ttlSeconds;
|
||||
|
||||
static final int DEFAULT_TTL_SECONDS = 60;
|
||||
|
||||
final Queue<ExecutorServiceExpiry> cache;
|
||||
final ThreadFactory factory;
|
||||
|
||||
final Queue<ExecutorService> all;
|
||||
final int ttlSeconds;
|
||||
|
||||
final ScheduledExecutorService evictor;
|
||||
|
||||
static final ExecutorService SHUTDOWN;
|
||||
static {
|
||||
SHUTDOWN = Executors.newSingleThreadExecutor();
|
||||
SHUTDOWN.shutdownNow();
|
||||
}
|
||||
|
||||
volatile boolean shutdown;
|
||||
|
||||
public ElasticScheduler(ThreadFactory factory, int ttlSeconds) {
|
||||
this.ttlSeconds = ttlSeconds;
|
||||
this.factory = factory;
|
||||
this.cache = new ConcurrentLinkedQueue<>();
|
||||
this.all = new ConcurrentLinkedQueue<>();
|
||||
this.evictor = Executors.newScheduledThreadPool(1, EVICTOR_FACTORY);
|
||||
this.evictor.scheduleAtFixedRate(this::eviction, ttlSeconds, ttlSeconds, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
throw new UnsupportedOperationException("Restarting not supported yet");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
shutdown = true;
|
||||
|
||||
evictor.shutdownNow();
|
||||
|
||||
cache.clear();
|
||||
|
||||
ExecutorService exec;
|
||||
|
||||
while ((exec = all.poll()) != null) {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
ExecutorService pick() {
|
||||
if (shutdown) {
|
||||
return SHUTDOWN;
|
||||
}
|
||||
ExecutorService result;
|
||||
ExecutorServiceExpiry e = cache.poll();
|
||||
if (e != null) {
|
||||
return e.executor;
|
||||
}
|
||||
|
||||
result = Executors.newSingleThreadExecutor(factory);
|
||||
all.offer(result);
|
||||
if (shutdown) {
|
||||
all.remove(result);
|
||||
return SHUTDOWN;
|
||||
}
|
||||
return result;
|
||||
}
|
||||
static final int DEFAULT_TTL_SECONDS = 60;
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
ExecutorService exec = pick();
|
||||
|
||||
Runnable wrapper = () -> {
|
||||
try {
|
||||
try {
|
||||
task.run();
|
||||
} catch (Throwable ex) {
|
||||
Exceptions.throwIfFatal(ex);
|
||||
Operators.onErrorDropped(ex);
|
||||
}
|
||||
} finally {
|
||||
release(exec);
|
||||
}
|
||||
};
|
||||
Future<?> f;
|
||||
|
||||
try {
|
||||
f = exec.submit(wrapper);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
return () -> f.cancel(true);
|
||||
}
|
||||
final Queue<ExecutorServiceExpiry> cache;
|
||||
|
||||
@Override
|
||||
public Worker createWorker() {
|
||||
ExecutorService exec = pick();
|
||||
return new CachedWorker(exec, this);
|
||||
}
|
||||
|
||||
void release(ExecutorService exec) {
|
||||
if (exec != SHUTDOWN && !shutdown) {
|
||||
ExecutorServiceExpiry e = new ExecutorServiceExpiry(exec, System.currentTimeMillis() + ttlSeconds * 1000L);
|
||||
cache.offer(e);
|
||||
if (shutdown) {
|
||||
if (cache.remove(e)) {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void eviction() {
|
||||
long now = System.currentTimeMillis();
|
||||
|
||||
List<ExecutorServiceExpiry> list = new ArrayList<>(cache);
|
||||
for (ExecutorServiceExpiry e : list) {
|
||||
if (e.expireMillis < now) {
|
||||
if (cache.remove(e)) {
|
||||
e.executor.shutdownNow();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
final Queue<ExecutorService> all;
|
||||
|
||||
static final class ExecutorServiceExpiry {
|
||||
final ExecutorService executor;
|
||||
final long expireMillis;
|
||||
final ScheduledExecutorService evictor;
|
||||
|
||||
public ExecutorServiceExpiry(ExecutorService executor, long expireMillis) {
|
||||
this.executor = executor;
|
||||
this.expireMillis = expireMillis;
|
||||
}
|
||||
}
|
||||
|
||||
static final class CachedWorker implements Worker {
|
||||
static final ExecutorService SHUTDOWN;
|
||||
|
||||
final ExecutorService executor;
|
||||
static {
|
||||
SHUTDOWN = Executors.newSingleThreadExecutor();
|
||||
SHUTDOWN.shutdownNow();
|
||||
}
|
||||
|
||||
final ElasticScheduler parent;
|
||||
volatile boolean shutdown;
|
||||
|
||||
volatile boolean shutdown;
|
||||
|
||||
OpenHashSet<CachedTask> tasks;
|
||||
|
||||
public CachedWorker(ExecutorService executor, ElasticScheduler parent) {
|
||||
this.executor = executor;
|
||||
this.parent = parent;
|
||||
this.tasks = new OpenHashSet<>();
|
||||
}
|
||||
public ElasticScheduler(ThreadFactory factory, int ttlSeconds) {
|
||||
this.ttlSeconds = ttlSeconds;
|
||||
this.factory = factory;
|
||||
this.cache = new ConcurrentLinkedQueue<>();
|
||||
this.all = new ConcurrentLinkedQueue<>();
|
||||
this.evictor = Executors.newScheduledThreadPool(1, EVICTOR_FACTORY);
|
||||
this.evictor.scheduleAtFixedRate(this::eviction, ttlSeconds, ttlSeconds, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
CachedTask ct = new CachedTask(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
tasks.add(ct);
|
||||
}
|
||||
|
||||
Future<?> f;
|
||||
try {
|
||||
f = executor.submit(ct);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
ct.setFuture(f);
|
||||
|
||||
return ct;
|
||||
}
|
||||
@Override
|
||||
public void start() {
|
||||
throw new UnsupportedOperationException("Restarting not supported yet");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
|
||||
OpenHashSet<CachedTask> set;
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
shutdown = true;
|
||||
set = tasks;
|
||||
tasks = null;
|
||||
}
|
||||
|
||||
if (!set.isEmpty()) {
|
||||
Object[] keys = set.keys();
|
||||
for (Object o : keys) {
|
||||
if (o != null) {
|
||||
((CachedTask)o).cancelFuture();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
parent.release(executor);
|
||||
}
|
||||
|
||||
void remove(CachedTask task) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
tasks.remove(task);
|
||||
}
|
||||
}
|
||||
|
||||
static final class CachedTask
|
||||
extends AtomicReference<Future<?>>
|
||||
implements Runnable, Cancellation {
|
||||
/** */
|
||||
private static final long serialVersionUID = 6799295393954430738L;
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
shutdown = true;
|
||||
|
||||
final Runnable run;
|
||||
|
||||
final CachedWorker parent;
|
||||
|
||||
volatile boolean cancelled;
|
||||
evictor.shutdownNow();
|
||||
|
||||
static final FutureTask<Object> CANCELLED = new FutureTask<>(() -> { }, null);
|
||||
cache.clear();
|
||||
|
||||
static final FutureTask<Object> FINISHED = new FutureTask<>(() -> { }, null);
|
||||
ExecutorService exec;
|
||||
|
||||
public CachedTask(Runnable run, CachedWorker parent) {
|
||||
this.run = run;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
try {
|
||||
if (!parent.shutdown && !cancelled) {
|
||||
run.run();
|
||||
}
|
||||
} catch (Throwable ex) {
|
||||
Exceptions.throwIfFatal(ex);
|
||||
Operators.onErrorDropped(ex);
|
||||
} finally {
|
||||
lazySet(FINISHED);
|
||||
parent.remove(this);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
cancelled = true;
|
||||
cancelFuture();
|
||||
}
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
if (!compareAndSet(null, f)) {
|
||||
if (get() != FINISHED) {
|
||||
f.cancel(true);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void cancelFuture() {
|
||||
Future<?> f = get();
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = getAndSet(CANCELLED);
|
||||
if (f != null && f != CANCELLED && f != FINISHED) {
|
||||
f.cancel(true);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
while ((exec = all.poll()) != null) {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
ExecutorService pick() {
|
||||
if (shutdown) {
|
||||
return SHUTDOWN;
|
||||
}
|
||||
ExecutorService result;
|
||||
ExecutorServiceExpiry e = cache.poll();
|
||||
if (e != null) {
|
||||
return e.executor;
|
||||
}
|
||||
|
||||
result = Executors.newSingleThreadExecutor(factory);
|
||||
all.offer(result);
|
||||
if (shutdown) {
|
||||
all.remove(result);
|
||||
return SHUTDOWN;
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
ExecutorService exec = pick();
|
||||
|
||||
Runnable wrapper = () -> {
|
||||
try {
|
||||
try {
|
||||
task.run();
|
||||
} catch (Throwable ex) {
|
||||
Exceptions.throwIfFatal(ex);
|
||||
Operators.onErrorDropped(ex);
|
||||
}
|
||||
} finally {
|
||||
release(exec);
|
||||
}
|
||||
};
|
||||
Future<?> f;
|
||||
|
||||
try {
|
||||
f = exec.submit(wrapper);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
return () -> f.cancel(true);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Worker createWorker() {
|
||||
ExecutorService exec = pick();
|
||||
return new CachedWorker(exec, this);
|
||||
}
|
||||
|
||||
void release(ExecutorService exec) {
|
||||
if (exec != SHUTDOWN && !shutdown) {
|
||||
ExecutorServiceExpiry e = new ExecutorServiceExpiry(exec, System.currentTimeMillis() + ttlSeconds * 1000L);
|
||||
cache.offer(e);
|
||||
if (shutdown) {
|
||||
if (cache.remove(e)) {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void eviction() {
|
||||
long now = System.currentTimeMillis();
|
||||
|
||||
List<ExecutorServiceExpiry> list = new ArrayList<>(cache);
|
||||
for (ExecutorServiceExpiry e : list) {
|
||||
if (e.expireMillis < now) {
|
||||
if (cache.remove(e)) {
|
||||
e.executor.shutdownNow();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
static final class ExecutorServiceExpiry {
|
||||
final ExecutorService executor;
|
||||
final long expireMillis;
|
||||
|
||||
public ExecutorServiceExpiry(ExecutorService executor, long expireMillis) {
|
||||
this.executor = executor;
|
||||
this.expireMillis = expireMillis;
|
||||
}
|
||||
}
|
||||
|
||||
static final class CachedWorker implements Worker {
|
||||
|
||||
final ExecutorService executor;
|
||||
|
||||
final ElasticScheduler parent;
|
||||
|
||||
volatile boolean shutdown;
|
||||
|
||||
OpenHashSet<CachedTask> tasks;
|
||||
|
||||
public CachedWorker(ExecutorService executor, ElasticScheduler parent) {
|
||||
this.executor = executor;
|
||||
this.parent = parent;
|
||||
this.tasks = new OpenHashSet<>();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
CachedTask ct = new CachedTask(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
tasks.add(ct);
|
||||
}
|
||||
|
||||
Future<?> f;
|
||||
try {
|
||||
f = executor.submit(ct);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
ct.setFuture(f);
|
||||
|
||||
return ct;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
|
||||
OpenHashSet<CachedTask> set;
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
shutdown = true;
|
||||
set = tasks;
|
||||
tasks = null;
|
||||
}
|
||||
|
||||
if (!set.isEmpty()) {
|
||||
Object[] keys = set.keys();
|
||||
for (Object o : keys) {
|
||||
if (o != null) {
|
||||
((CachedTask) o).cancelFuture();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
parent.release(executor);
|
||||
}
|
||||
|
||||
void remove(CachedTask task) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
tasks.remove(task);
|
||||
}
|
||||
}
|
||||
|
||||
static final class CachedTask
|
||||
extends AtomicReference<Future<?>>
|
||||
implements Runnable, Cancellation {
|
||||
/** */
|
||||
private static final long serialVersionUID = 6799295393954430738L;
|
||||
|
||||
final Runnable run;
|
||||
|
||||
final CachedWorker parent;
|
||||
|
||||
volatile boolean cancelled;
|
||||
|
||||
static final FutureTask<Object> CANCELLED = new FutureTask<>(() -> {
|
||||
}, null);
|
||||
|
||||
static final FutureTask<Object> FINISHED = new FutureTask<>(() -> {
|
||||
}, null);
|
||||
|
||||
public CachedTask(Runnable run, CachedWorker parent) {
|
||||
this.run = run;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
try {
|
||||
if (!parent.shutdown && !cancelled) {
|
||||
run.run();
|
||||
}
|
||||
} catch (Throwable ex) {
|
||||
Exceptions.throwIfFatal(ex);
|
||||
Operators.onErrorDropped(ex);
|
||||
} finally {
|
||||
lazySet(FINISHED);
|
||||
parent.remove(this);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
cancelled = true;
|
||||
cancelFuture();
|
||||
}
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
if (!compareAndSet(null, f)) {
|
||||
if (get() != FINISHED) {
|
||||
f.cancel(true);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void cancelFuture() {
|
||||
Future<?> f = get();
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = getAndSet(CANCELLED);
|
||||
if (f != null && f != CANCELLED && f != FINISHED) {
|
||||
f.cancel(true);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -33,291 +33,293 @@ import reactor.util.concurrent.OpenHashSet;
|
||||
/**
|
||||
* Scheduler that hosts a fixed pool of single-threaded ExecutorService-based workers
|
||||
* and is suited for parallel work.
|
||||
*
|
||||
* @author Stephane Maldini
|
||||
*/
|
||||
final class ParallelScheduler implements Scheduler {
|
||||
|
||||
static final AtomicLong COUNTER = new AtomicLong();
|
||||
static final AtomicLong COUNTER = new AtomicLong();
|
||||
|
||||
final int n;
|
||||
|
||||
final ThreadFactory factory;
|
||||
final int n;
|
||||
|
||||
volatile ExecutorService[] executors;
|
||||
static final AtomicReferenceFieldUpdater<ParallelScheduler, ExecutorService[]> EXECUTORS =
|
||||
AtomicReferenceFieldUpdater.newUpdater(ParallelScheduler.class, ExecutorService[].class, "executors");
|
||||
final ThreadFactory factory;
|
||||
|
||||
static final ExecutorService[] SHUTDOWN = new ExecutorService[0];
|
||||
|
||||
static final ExecutorService TERMINATED;
|
||||
static {
|
||||
TERMINATED = Executors.newSingleThreadExecutor();
|
||||
TERMINATED.shutdownNow();
|
||||
}
|
||||
|
||||
int roundRobin;
|
||||
volatile ExecutorService[] executors;
|
||||
static final AtomicReferenceFieldUpdater<ParallelScheduler, ExecutorService[]> EXECUTORS =
|
||||
AtomicReferenceFieldUpdater.newUpdater(ParallelScheduler.class, ExecutorService[].class, "executors");
|
||||
|
||||
ParallelScheduler(int n, ThreadFactory factory) {
|
||||
if (n <= 0) {
|
||||
throw new IllegalArgumentException("n > 0 required but it was " + n);
|
||||
}
|
||||
this.n = n;
|
||||
this.factory = factory;
|
||||
init(n);
|
||||
}
|
||||
|
||||
void init(int n) {
|
||||
ExecutorService[] a = new ExecutorService[n];
|
||||
for (int i = 0; i < n; i++) {
|
||||
a[i] = Executors.newSingleThreadExecutor(factory);
|
||||
}
|
||||
EXECUTORS.lazySet(this, a);
|
||||
}
|
||||
static final ExecutorService[] SHUTDOWN = new ExecutorService[0];
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
ExecutorService[] b = null;
|
||||
for (;;) {
|
||||
ExecutorService[] a = executors;
|
||||
if (a != SHUTDOWN) {
|
||||
if (b != null) {
|
||||
for (ExecutorService exec : b) {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
static final ExecutorService TERMINATED;
|
||||
|
||||
if (b == null) {
|
||||
b = new ExecutorService[n];
|
||||
for (int i = 0; i < n; i++) {
|
||||
b[i] = Executors.newSingleThreadExecutor(factory);
|
||||
}
|
||||
}
|
||||
|
||||
if (EXECUTORS.compareAndSet(this, a, b)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
ExecutorService[] a = executors;
|
||||
if (a != SHUTDOWN) {
|
||||
a = EXECUTORS.getAndSet(this, SHUTDOWN);
|
||||
if (a != SHUTDOWN) {
|
||||
for (ExecutorService exec : a) {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
ExecutorService pick() {
|
||||
ExecutorService[] a = executors;
|
||||
if (a != SHUTDOWN) {
|
||||
// ignoring the race condition here, its already random who gets which executor
|
||||
int idx = roundRobin;
|
||||
if (idx == n) {
|
||||
idx = 0;
|
||||
roundRobin = 0;
|
||||
} else {
|
||||
roundRobin = idx + 1;
|
||||
}
|
||||
return a[idx];
|
||||
}
|
||||
return TERMINATED;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
ExecutorService exec = pick();
|
||||
Future<?> f = exec.submit(task);
|
||||
return () -> f.cancel(false);
|
||||
}
|
||||
static {
|
||||
TERMINATED = Executors.newSingleThreadExecutor();
|
||||
TERMINATED.shutdownNow();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Worker createWorker() {
|
||||
return new ParallelWorker(pick());
|
||||
}
|
||||
|
||||
static final class ParallelWorker implements Worker {
|
||||
final ExecutorService exec;
|
||||
|
||||
OpenHashSet<ParallelWorkerTask> tasks;
|
||||
|
||||
volatile boolean shutdown;
|
||||
|
||||
public ParallelWorker(ExecutorService exec) {
|
||||
this.exec = exec;
|
||||
this.tasks = new OpenHashSet<>();
|
||||
}
|
||||
int roundRobin;
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
ParallelWorkerTask pw = new ParallelWorkerTask(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
tasks.add(pw);
|
||||
}
|
||||
|
||||
Future<?> f;
|
||||
try {
|
||||
f = exec.submit(pw);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
if (shutdown) {
|
||||
f.cancel(true);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
pw.setFuture(f);
|
||||
|
||||
return pw;
|
||||
}
|
||||
ParallelScheduler(int n, ThreadFactory factory) {
|
||||
if (n <= 0) {
|
||||
throw new IllegalArgumentException("n > 0 required but it was " + n);
|
||||
}
|
||||
this.n = n;
|
||||
this.factory = factory;
|
||||
init(n);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
shutdown = true;
|
||||
OpenHashSet<ParallelWorkerTask> set;
|
||||
synchronized (this) {
|
||||
set = tasks;
|
||||
tasks = null;
|
||||
}
|
||||
|
||||
if (set != null) {
|
||||
Object[] a = set.keys();
|
||||
for (Object o : a) {
|
||||
if (o != null) {
|
||||
((ParallelWorkerTask)o).cancelFuture();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void remove(ParallelWorkerTask task) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
tasks.remove(task);
|
||||
}
|
||||
}
|
||||
|
||||
int pendingTasks() {
|
||||
if (shutdown) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
OpenHashSet<?> set = tasks;
|
||||
if (set != null) {
|
||||
return set.size();
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
static final class ParallelWorkerTask implements Runnable, Cancellation {
|
||||
final Runnable run;
|
||||
|
||||
final ParallelWorker parent;
|
||||
|
||||
volatile boolean cancelled;
|
||||
|
||||
volatile Future<?> future;
|
||||
@SuppressWarnings("rawtypes")
|
||||
static final AtomicReferenceFieldUpdater<ParallelWorkerTask, Future> FUTURE =
|
||||
AtomicReferenceFieldUpdater.newUpdater(ParallelWorkerTask.class, Future.class, "future");
|
||||
|
||||
static final Future<Object> FINISHED = CompletableFuture.completedFuture(null);
|
||||
static final Future<Object> CANCELLED = CompletableFuture.completedFuture(null);
|
||||
|
||||
public ParallelWorkerTask(Runnable run, ParallelWorker parent) {
|
||||
this.run = run;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
if (cancelled || parent.shutdown) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
try {
|
||||
run.run();
|
||||
} catch (Throwable ex) {
|
||||
Exceptions.throwIfFatal(ex);
|
||||
Operators.onErrorDropped(ex);
|
||||
}
|
||||
} finally {
|
||||
for (;;) {
|
||||
Future<?> f = future;
|
||||
if (f == CANCELLED) {
|
||||
break;
|
||||
}
|
||||
if (FUTURE.compareAndSet(this, f, FINISHED)) {
|
||||
parent.remove(this);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
if (!cancelled) {
|
||||
cancelled = true;
|
||||
|
||||
Future<?> f = future;
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = FUTURE.getAndSet(this, CANCELLED);
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
if (f != null) {
|
||||
f.cancel(parent.shutdown);
|
||||
}
|
||||
|
||||
parent.remove(this);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
if (future != null || !FUTURE.compareAndSet(this, null, f)) {
|
||||
if (future != FINISHED) {
|
||||
f.cancel(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void cancelFuture() {
|
||||
Future<?> f = future;
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = FUTURE.getAndSet(this, CANCELLED);
|
||||
if (f != null && f != CANCELLED && f != FINISHED) {
|
||||
f.cancel(true);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
void init(int n) {
|
||||
ExecutorService[] a = new ExecutorService[n];
|
||||
for (int i = 0; i < n; i++) {
|
||||
a[i] = Executors.newSingleThreadExecutor(factory);
|
||||
}
|
||||
EXECUTORS.lazySet(this, a);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
ExecutorService[] b = null;
|
||||
for (; ; ) {
|
||||
ExecutorService[] a = executors;
|
||||
if (a != SHUTDOWN) {
|
||||
if (b != null) {
|
||||
for (ExecutorService exec : b) {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (b == null) {
|
||||
b = new ExecutorService[n];
|
||||
for (int i = 0; i < n; i++) {
|
||||
b[i] = Executors.newSingleThreadExecutor(factory);
|
||||
}
|
||||
}
|
||||
|
||||
if (EXECUTORS.compareAndSet(this, a, b)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
ExecutorService[] a = executors;
|
||||
if (a != SHUTDOWN) {
|
||||
a = EXECUTORS.getAndSet(this, SHUTDOWN);
|
||||
if (a != SHUTDOWN) {
|
||||
for (ExecutorService exec : a) {
|
||||
exec.shutdownNow();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
ExecutorService pick() {
|
||||
ExecutorService[] a = executors;
|
||||
if (a != SHUTDOWN) {
|
||||
// ignoring the race condition here, its already random who gets which executor
|
||||
int idx = roundRobin;
|
||||
if (idx == n) {
|
||||
idx = 0;
|
||||
roundRobin = 0;
|
||||
} else {
|
||||
roundRobin = idx + 1;
|
||||
}
|
||||
return a[idx];
|
||||
}
|
||||
return TERMINATED;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
ExecutorService exec = pick();
|
||||
Future<?> f = exec.submit(task);
|
||||
return () -> f.cancel(false);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Worker createWorker() {
|
||||
return new ParallelWorker(pick());
|
||||
}
|
||||
|
||||
static final class ParallelWorker implements Worker {
|
||||
final ExecutorService exec;
|
||||
|
||||
OpenHashSet<ParallelWorkerTask> tasks;
|
||||
|
||||
volatile boolean shutdown;
|
||||
|
||||
public ParallelWorker(ExecutorService exec) {
|
||||
this.exec = exec;
|
||||
this.tasks = new OpenHashSet<>();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
ParallelWorkerTask pw = new ParallelWorkerTask(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
tasks.add(pw);
|
||||
}
|
||||
|
||||
Future<?> f;
|
||||
try {
|
||||
f = exec.submit(pw);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
if (shutdown) {
|
||||
f.cancel(true);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
pw.setFuture(f);
|
||||
|
||||
return pw;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
shutdown = true;
|
||||
OpenHashSet<ParallelWorkerTask> set;
|
||||
synchronized (this) {
|
||||
set = tasks;
|
||||
tasks = null;
|
||||
}
|
||||
|
||||
if (set != null) {
|
||||
Object[] a = set.keys();
|
||||
for (Object o : a) {
|
||||
if (o != null) {
|
||||
((ParallelWorkerTask) o).cancelFuture();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void remove(ParallelWorkerTask task) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
tasks.remove(task);
|
||||
}
|
||||
}
|
||||
|
||||
int pendingTasks() {
|
||||
if (shutdown) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
OpenHashSet<?> set = tasks;
|
||||
if (set != null) {
|
||||
return set.size();
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
static final class ParallelWorkerTask implements Runnable, Cancellation {
|
||||
final Runnable run;
|
||||
|
||||
final ParallelWorker parent;
|
||||
|
||||
volatile boolean cancelled;
|
||||
|
||||
volatile Future<?> future;
|
||||
@SuppressWarnings("rawtypes")
|
||||
static final AtomicReferenceFieldUpdater<ParallelWorkerTask, Future> FUTURE =
|
||||
AtomicReferenceFieldUpdater.newUpdater(ParallelWorkerTask.class, Future.class, "future");
|
||||
|
||||
static final Future<Object> FINISHED = CompletableFuture.completedFuture(null);
|
||||
static final Future<Object> CANCELLED = CompletableFuture.completedFuture(null);
|
||||
|
||||
public ParallelWorkerTask(Runnable run, ParallelWorker parent) {
|
||||
this.run = run;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
if (cancelled || parent.shutdown) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
try {
|
||||
run.run();
|
||||
} catch (Throwable ex) {
|
||||
Exceptions.throwIfFatal(ex);
|
||||
Operators.onErrorDropped(ex);
|
||||
}
|
||||
} finally {
|
||||
for (; ; ) {
|
||||
Future<?> f = future;
|
||||
if (f == CANCELLED) {
|
||||
break;
|
||||
}
|
||||
if (FUTURE.compareAndSet(this, f, FINISHED)) {
|
||||
parent.remove(this);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
if (!cancelled) {
|
||||
cancelled = true;
|
||||
|
||||
Future<?> f = future;
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = FUTURE.getAndSet(this, CANCELLED);
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
if (f != null) {
|
||||
f.cancel(parent.shutdown);
|
||||
}
|
||||
|
||||
parent.remove(this);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
if (future != null || !FUTURE.compareAndSet(this, null, f)) {
|
||||
if (future != FINISHED) {
|
||||
f.cancel(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void cancelFuture() {
|
||||
Future<?> f = future;
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = FUTURE.getAndSet(this, CANCELLED);
|
||||
if (f != null && f != CANCELLED && f != FINISHED) {
|
||||
f.cancel(true);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -33,261 +33,263 @@ import reactor.util.concurrent.OpenHashSet;
|
||||
/**
|
||||
* Scheduler that works with a single-threaded ExecutorService and is suited for
|
||||
* same-thread work (like an event dispatch thread).
|
||||
* @author Stephane Maldini
|
||||
*/
|
||||
final class SingleScheduler implements Scheduler {
|
||||
|
||||
static final AtomicLong COUNTER = new AtomicLong();
|
||||
|
||||
final ThreadFactory factory;
|
||||
static final AtomicLong COUNTER = new AtomicLong();
|
||||
|
||||
volatile ExecutorService executor;
|
||||
static final AtomicReferenceFieldUpdater<SingleScheduler, ExecutorService> EXECUTORS =
|
||||
AtomicReferenceFieldUpdater.newUpdater(SingleScheduler.class, ExecutorService.class, "executor");
|
||||
final ThreadFactory factory;
|
||||
|
||||
static final ExecutorService TERMINATED;
|
||||
static {
|
||||
TERMINATED = Executors.newSingleThreadExecutor();
|
||||
TERMINATED.shutdownNow();
|
||||
}
|
||||
|
||||
public SingleScheduler(ThreadFactory factory) {
|
||||
this.factory = factory;
|
||||
init();
|
||||
}
|
||||
|
||||
private void init() {
|
||||
EXECUTORS.lazySet(this, Executors.newSingleThreadExecutor(factory));
|
||||
}
|
||||
|
||||
public boolean isStarted() {
|
||||
return executor != TERMINATED;
|
||||
}
|
||||
volatile ExecutorService executor;
|
||||
static final AtomicReferenceFieldUpdater<SingleScheduler, ExecutorService> EXECUTORS =
|
||||
AtomicReferenceFieldUpdater.newUpdater(SingleScheduler.class, ExecutorService.class, "executor");
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
ExecutorService b = null;
|
||||
for (;;) {
|
||||
ExecutorService a = executor;
|
||||
if (a != TERMINATED) {
|
||||
if (b != null) {
|
||||
b.shutdownNow();
|
||||
}
|
||||
return;
|
||||
}
|
||||
static final ExecutorService TERMINATED;
|
||||
|
||||
if (b == null) {
|
||||
b = Executors.newSingleThreadExecutor(factory);
|
||||
}
|
||||
|
||||
if (EXECUTORS.compareAndSet(this, a, b)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
ExecutorService a = executor;
|
||||
if (a != TERMINATED) {
|
||||
a = EXECUTORS.getAndSet(this, TERMINATED);
|
||||
if (a != TERMINATED) {
|
||||
a.shutdownNow();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
try {
|
||||
Future<?> f = executor.submit(task);
|
||||
return () -> f.cancel(false);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
}
|
||||
static {
|
||||
TERMINATED = Executors.newSingleThreadExecutor();
|
||||
TERMINATED.shutdownNow();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Worker createWorker() {
|
||||
return new SingleWorker(executor);
|
||||
}
|
||||
|
||||
static final class SingleWorker implements Worker {
|
||||
final ExecutorService exec;
|
||||
|
||||
OpenHashSet<SingleWorkerTask> tasks;
|
||||
|
||||
volatile boolean shutdown;
|
||||
|
||||
public SingleWorker(ExecutorService exec) {
|
||||
this.exec = exec;
|
||||
this.tasks = new OpenHashSet<>();
|
||||
}
|
||||
public SingleScheduler(ThreadFactory factory) {
|
||||
this.factory = factory;
|
||||
init();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
SingleWorkerTask pw = new SingleWorkerTask(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
tasks.add(pw);
|
||||
}
|
||||
|
||||
Future<?> f;
|
||||
try {
|
||||
f = exec.submit(pw);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
if (shutdown) {
|
||||
f.cancel(true);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
pw.setFuture(f);
|
||||
|
||||
return pw;
|
||||
}
|
||||
private void init() {
|
||||
EXECUTORS.lazySet(this, Executors.newSingleThreadExecutor(factory));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
shutdown = true;
|
||||
OpenHashSet<SingleWorkerTask> set;
|
||||
synchronized (this) {
|
||||
set = tasks;
|
||||
tasks = null;
|
||||
}
|
||||
|
||||
if (set != null && !set.isEmpty()) {
|
||||
Object[] a = set.keys();
|
||||
for (Object o : a) {
|
||||
if (o != null) {
|
||||
((SingleWorkerTask)o).cancelFuture();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void remove(SingleWorkerTask task) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
tasks.remove(task);
|
||||
}
|
||||
}
|
||||
|
||||
int pendingTasks() {
|
||||
if (shutdown) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
OpenHashSet<?> set = tasks;
|
||||
if (set != null) {
|
||||
return set.size();
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
static final class SingleWorkerTask implements Runnable, Cancellation {
|
||||
final Runnable run;
|
||||
|
||||
final SingleWorker parent;
|
||||
|
||||
volatile boolean cancelled;
|
||||
|
||||
volatile Future<?> future;
|
||||
@SuppressWarnings("rawtypes")
|
||||
static final AtomicReferenceFieldUpdater<SingleWorkerTask, Future> FUTURE =
|
||||
AtomicReferenceFieldUpdater.newUpdater(SingleWorkerTask.class, Future.class, "future");
|
||||
|
||||
static final Future<Object> FINISHED = CompletableFuture.completedFuture(null);
|
||||
static final Future<Object> CANCELLED = CompletableFuture.completedFuture(null);
|
||||
|
||||
public SingleWorkerTask(Runnable run, SingleWorker parent) {
|
||||
this.run = run;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
if (cancelled || parent.shutdown) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
try {
|
||||
run.run();
|
||||
} catch (Throwable ex) {
|
||||
Exceptions.throwIfFatal(ex);
|
||||
Operators.onErrorDropped(ex);
|
||||
}
|
||||
} finally {
|
||||
for (;;) {
|
||||
Future<?> f = future;
|
||||
if (f == CANCELLED) {
|
||||
break;
|
||||
}
|
||||
if (FUTURE.compareAndSet(this, f, FINISHED)) {
|
||||
parent.remove(this);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
if (!cancelled) {
|
||||
cancelled = true;
|
||||
|
||||
Future<?> f = future;
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = FUTURE.getAndSet(this, CANCELLED);
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
if (f != null) {
|
||||
f.cancel(false);
|
||||
}
|
||||
|
||||
parent.remove(this);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
if (future != null || !FUTURE.compareAndSet(this, null, f)) {
|
||||
if (future != FINISHED) {
|
||||
f.cancel(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void cancelFuture() {
|
||||
Future<?> f = future;
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = FUTURE.getAndSet(this, CANCELLED);
|
||||
if (f != null && f != CANCELLED && f != FINISHED) {
|
||||
f.cancel(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
public boolean isStarted() {
|
||||
return executor != TERMINATED;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
ExecutorService b = null;
|
||||
for (; ; ) {
|
||||
ExecutorService a = executor;
|
||||
if (a != TERMINATED) {
|
||||
if (b != null) {
|
||||
b.shutdownNow();
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (b == null) {
|
||||
b = Executors.newSingleThreadExecutor(factory);
|
||||
}
|
||||
|
||||
if (EXECUTORS.compareAndSet(this, a, b)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
ExecutorService a = executor;
|
||||
if (a != TERMINATED) {
|
||||
a = EXECUTORS.getAndSet(this, TERMINATED);
|
||||
if (a != TERMINATED) {
|
||||
a.shutdownNow();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
try {
|
||||
Future<?> f = executor.submit(task);
|
||||
return () -> f.cancel(false);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Worker createWorker() {
|
||||
return new SingleWorker(executor);
|
||||
}
|
||||
|
||||
static final class SingleWorker implements Worker {
|
||||
final ExecutorService exec;
|
||||
|
||||
OpenHashSet<SingleWorkerTask> tasks;
|
||||
|
||||
volatile boolean shutdown;
|
||||
|
||||
public SingleWorker(ExecutorService exec) {
|
||||
this.exec = exec;
|
||||
this.tasks = new OpenHashSet<>();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
SingleWorkerTask pw = new SingleWorkerTask(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return REJECTED;
|
||||
}
|
||||
tasks.add(pw);
|
||||
}
|
||||
|
||||
Future<?> f;
|
||||
try {
|
||||
f = exec.submit(pw);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
if (shutdown) {
|
||||
f.cancel(true);
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
pw.setFuture(f);
|
||||
|
||||
return pw;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
shutdown = true;
|
||||
OpenHashSet<SingleWorkerTask> set;
|
||||
synchronized (this) {
|
||||
set = tasks;
|
||||
tasks = null;
|
||||
}
|
||||
|
||||
if (set != null && !set.isEmpty()) {
|
||||
Object[] a = set.keys();
|
||||
for (Object o : a) {
|
||||
if (o != null) {
|
||||
((SingleWorkerTask) o).cancelFuture();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void remove(SingleWorkerTask task) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
if (shutdown) {
|
||||
return;
|
||||
}
|
||||
tasks.remove(task);
|
||||
}
|
||||
}
|
||||
|
||||
int pendingTasks() {
|
||||
if (shutdown) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
synchronized (this) {
|
||||
OpenHashSet<?> set = tasks;
|
||||
if (set != null) {
|
||||
return set.size();
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
static final class SingleWorkerTask implements Runnable, Cancellation {
|
||||
final Runnable run;
|
||||
|
||||
final SingleWorker parent;
|
||||
|
||||
volatile boolean cancelled;
|
||||
|
||||
volatile Future<?> future;
|
||||
@SuppressWarnings("rawtypes")
|
||||
static final AtomicReferenceFieldUpdater<SingleWorkerTask, Future> FUTURE =
|
||||
AtomicReferenceFieldUpdater.newUpdater(SingleWorkerTask.class, Future.class, "future");
|
||||
|
||||
static final Future<Object> FINISHED = CompletableFuture.completedFuture(null);
|
||||
static final Future<Object> CANCELLED = CompletableFuture.completedFuture(null);
|
||||
|
||||
public SingleWorkerTask(Runnable run, SingleWorker parent) {
|
||||
this.run = run;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
if (cancelled || parent.shutdown) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
try {
|
||||
run.run();
|
||||
} catch (Throwable ex) {
|
||||
Exceptions.throwIfFatal(ex);
|
||||
Operators.onErrorDropped(ex);
|
||||
}
|
||||
} finally {
|
||||
for (; ; ) {
|
||||
Future<?> f = future;
|
||||
if (f == CANCELLED) {
|
||||
break;
|
||||
}
|
||||
if (FUTURE.compareAndSet(this, f, FINISHED)) {
|
||||
parent.remove(this);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
if (!cancelled) {
|
||||
cancelled = true;
|
||||
|
||||
Future<?> f = future;
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = FUTURE.getAndSet(this, CANCELLED);
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
if (f != null) {
|
||||
f.cancel(false);
|
||||
}
|
||||
|
||||
parent.remove(this);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
if (future != null || !FUTURE.compareAndSet(this, null, f)) {
|
||||
if (future != FINISHED) {
|
||||
f.cancel(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void cancelFuture() {
|
||||
Future<?> f = future;
|
||||
if (f != CANCELLED && f != FINISHED) {
|
||||
f = FUTURE.getAndSet(this, CANCELLED);
|
||||
if (f != null && f != CANCELLED && f != FINISHED) {
|
||||
f.cancel(false);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -37,428 +37,429 @@ import reactor.util.concurrent.OpenHashSet;
|
||||
*/
|
||||
final class SingleTimedScheduler implements TimedScheduler {
|
||||
|
||||
static final AtomicLong COUNTER = new AtomicLong();
|
||||
|
||||
final ScheduledThreadPoolExecutor executor;
|
||||
static final AtomicLong COUNTER = new AtomicLong();
|
||||
|
||||
/**
|
||||
* Constructs a new SingleTimedScheduler with the given thread factory.
|
||||
* @param threadFactory the thread factory to use
|
||||
*/
|
||||
SingleTimedScheduler(ThreadFactory threadFactory) {
|
||||
ScheduledThreadPoolExecutor e = (ScheduledThreadPoolExecutor)Executors.newScheduledThreadPool(1, threadFactory);
|
||||
e.setRemoveOnCancelPolicy(true);
|
||||
executor = e;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
try {
|
||||
Future<?> f = executor.submit(task);
|
||||
return () -> f.cancel(false);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
return REJECTED;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task, long delay, TimeUnit unit) {
|
||||
try {
|
||||
Future<?> f = executor.schedule(task, delay, unit);
|
||||
return () -> f.cancel(false);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
return REJECTED;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedulePeriodically(Runnable task, long initialDelay, long period, TimeUnit unit) {
|
||||
try {
|
||||
Future<?> f = executor.scheduleAtFixedRate(task, initialDelay, period, unit);
|
||||
return () -> f.cancel(false);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
return REJECTED;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
throw new UnsupportedOperationException("Not supported, yet.");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
executor.shutdownNow();
|
||||
}
|
||||
|
||||
@Override
|
||||
public TimedWorker createWorker() {
|
||||
return new SingleTimedSchedulerWorker(executor);
|
||||
}
|
||||
|
||||
static final class SingleTimedSchedulerWorker implements TimedWorker {
|
||||
final ScheduledThreadPoolExecutor executor;
|
||||
|
||||
OpenHashSet<CancelFuture> tasks;
|
||||
|
||||
volatile boolean terminated;
|
||||
|
||||
public SingleTimedSchedulerWorker(ScheduledThreadPoolExecutor executor) {
|
||||
this.executor = executor;
|
||||
this.tasks = new OpenHashSet<>();
|
||||
}
|
||||
final ScheduledThreadPoolExecutor executor;
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
TimedScheduledRunnable sr = new TimedScheduledRunnable(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
tasks.add(sr);
|
||||
}
|
||||
|
||||
try {
|
||||
Future<?> f = executor.submit(sr);
|
||||
sr.set(f);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
sr.dispose();
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
return sr;
|
||||
}
|
||||
|
||||
void delete(CancelFuture r) {
|
||||
synchronized (this) {
|
||||
if (!terminated) {
|
||||
tasks.remove(r);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task, long delay, TimeUnit unit) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
TimedScheduledRunnable sr = new TimedScheduledRunnable(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
tasks.add(sr);
|
||||
}
|
||||
|
||||
try {
|
||||
Future<?> f = executor.schedule(sr, delay, unit);
|
||||
sr.set(f);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
sr.dispose();
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
return sr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedulePeriodically(Runnable task, long initialDelay, long period, TimeUnit unit) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
TimedPeriodicScheduledRunnable sr = new TimedPeriodicScheduledRunnable(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
tasks.add(sr);
|
||||
}
|
||||
|
||||
try {
|
||||
Future<?> f = executor.scheduleAtFixedRate(sr, initialDelay, period, unit);
|
||||
sr.set(f);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
sr.dispose();
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
return sr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (terminated) {
|
||||
return;
|
||||
}
|
||||
terminated = true;
|
||||
|
||||
OpenHashSet<CancelFuture> set;
|
||||
|
||||
synchronized (this) {
|
||||
set = tasks;
|
||||
if (set == null) {
|
||||
return;
|
||||
}
|
||||
tasks = null;
|
||||
}
|
||||
|
||||
if (!set.isEmpty()) {
|
||||
Object[] keys = set.keys();
|
||||
for (Object c : keys) {
|
||||
if (c != null) {
|
||||
((CancelFuture)c).cancelFuture();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
/**
|
||||
* Constructs a new SingleTimedScheduler with the given thread factory.
|
||||
*
|
||||
* @param threadFactory the thread factory to use
|
||||
*/
|
||||
SingleTimedScheduler(ThreadFactory threadFactory) {
|
||||
ScheduledThreadPoolExecutor e = (ScheduledThreadPoolExecutor) Executors.newScheduledThreadPool(1, threadFactory);
|
||||
e.setRemoveOnCancelPolicy(true);
|
||||
executor = e;
|
||||
}
|
||||
|
||||
interface CancelFuture {
|
||||
void cancelFuture();
|
||||
}
|
||||
|
||||
static final class TimedScheduledRunnable
|
||||
extends AtomicReference<Future<?>> implements Runnable, Cancellation, CancelFuture {
|
||||
/** */
|
||||
private static final long serialVersionUID = 2284024836904862408L;
|
||||
|
||||
final Runnable task;
|
||||
|
||||
final SingleTimedSchedulerWorker parent;
|
||||
|
||||
volatile Thread current;
|
||||
static final AtomicReferenceFieldUpdater<TimedScheduledRunnable, Thread> CURRENT =
|
||||
AtomicReferenceFieldUpdater.newUpdater(TimedScheduledRunnable.class, Thread.class, "current");
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
try {
|
||||
Future<?> f = executor.submit(task);
|
||||
return () -> f.cancel(false);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
return REJECTED;
|
||||
}
|
||||
}
|
||||
|
||||
static final Runnable EMPTY = new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task, long delay, TimeUnit unit) {
|
||||
try {
|
||||
Future<?> f = executor.schedule(task, delay, unit);
|
||||
return () -> f.cancel(false);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
return REJECTED;
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
};
|
||||
@Override
|
||||
public Cancellation schedulePeriodically(Runnable task, long initialDelay, long period, TimeUnit unit) {
|
||||
try {
|
||||
Future<?> f = executor.scheduleAtFixedRate(task, initialDelay, period, unit);
|
||||
return () -> f.cancel(false);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
return REJECTED;
|
||||
}
|
||||
}
|
||||
|
||||
static final Future<?> CANCELLED_FUTURE = new FutureTask<>(EMPTY, null);
|
||||
@Override
|
||||
public void start() {
|
||||
throw new UnsupportedOperationException("Not supported, yet.");
|
||||
}
|
||||
|
||||
static final Future<?> FINISHED = new FutureTask<>(EMPTY, null);
|
||||
@Override
|
||||
public void shutdown() {
|
||||
executor.shutdownNow();
|
||||
}
|
||||
|
||||
public TimedScheduledRunnable(Runnable task, SingleTimedSchedulerWorker parent) {
|
||||
this.task = task;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
CURRENT.lazySet(this, Thread.currentThread());
|
||||
try {
|
||||
try {
|
||||
task.run();
|
||||
} catch (Throwable e) {
|
||||
Operators.onErrorDropped(e);
|
||||
}
|
||||
} finally {
|
||||
for (;;) {
|
||||
Future<?> a = get();
|
||||
if (a == CANCELLED_FUTURE) {
|
||||
break;
|
||||
}
|
||||
if (compareAndSet(a, FINISHED)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
parent.delete(this);
|
||||
break;
|
||||
}
|
||||
}
|
||||
CURRENT.lazySet(this, null);
|
||||
}
|
||||
}
|
||||
|
||||
void doCancel(Future<?> a) {
|
||||
a.cancel(Thread.currentThread() != current);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void cancelFuture() {
|
||||
for (;;) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(a, CANCELLED_FUTURE)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
@Override
|
||||
public TimedWorker createWorker() {
|
||||
return new SingleTimedSchedulerWorker(executor);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
for (;;) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(a, CANCELLED_FUTURE)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
parent.delete(this);
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
static final class SingleTimedSchedulerWorker implements TimedWorker {
|
||||
final ScheduledThreadPoolExecutor executor;
|
||||
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
for (;;) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (a == CANCELLED_FUTURE) {
|
||||
doCancel(a);
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(null, f)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "TimedScheduledRunnable[cancelled=" + (get() == CANCELLED_FUTURE) +
|
||||
", task=" + task +
|
||||
"]";
|
||||
}
|
||||
}
|
||||
OpenHashSet<CancelFuture> tasks;
|
||||
|
||||
static final class TimedPeriodicScheduledRunnable
|
||||
extends AtomicReference<Future<?>>
|
||||
implements Runnable, Cancellation, CancelFuture {
|
||||
/** */
|
||||
private static final long serialVersionUID = 2284024836904862408L;
|
||||
|
||||
final Runnable task;
|
||||
|
||||
final SingleTimedSchedulerWorker parent;
|
||||
|
||||
volatile Thread current;
|
||||
static final AtomicReferenceFieldUpdater<TimedPeriodicScheduledRunnable, Thread> CURRENT =
|
||||
AtomicReferenceFieldUpdater.newUpdater(TimedPeriodicScheduledRunnable.class, Thread.class, "current");
|
||||
volatile boolean terminated;
|
||||
|
||||
static final Runnable EMPTY = new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
public SingleTimedSchedulerWorker(ScheduledThreadPoolExecutor executor) {
|
||||
this.executor = executor;
|
||||
this.tasks = new OpenHashSet<>();
|
||||
}
|
||||
|
||||
}
|
||||
};
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
static final Future<?> CANCELLED_FUTURE = new FutureTask<>(EMPTY, null);
|
||||
TimedScheduledRunnable sr = new TimedScheduledRunnable(task, this);
|
||||
|
||||
static final Future<?> FINISHED = new FutureTask<>(EMPTY, null);
|
||||
synchronized (this) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
public TimedPeriodicScheduledRunnable(Runnable task, SingleTimedSchedulerWorker parent) {
|
||||
this.task = task;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
CURRENT.lazySet(this, Thread.currentThread());
|
||||
try {
|
||||
try {
|
||||
task.run();
|
||||
} catch (Throwable ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
for (;;) {
|
||||
Future<?> a = get();
|
||||
if (a == CANCELLED_FUTURE) {
|
||||
break;
|
||||
}
|
||||
if (compareAndSet(a, FINISHED)) {
|
||||
parent.delete(this);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
CURRENT.lazySet(this, null);
|
||||
}
|
||||
}
|
||||
|
||||
void doCancel(Future<?> a) {
|
||||
a.cancel(false);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void cancelFuture() {
|
||||
for (;;) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(a, CANCELLED_FUTURE)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
for (;;) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(a, CANCELLED_FUTURE)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
parent.delete(this);
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
tasks.add(sr);
|
||||
}
|
||||
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
for (;;) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (a == CANCELLED_FUTURE) {
|
||||
doCancel(a);
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(null, f)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "TimedPeriodicScheduledRunnable[cancelled=" + get() + ", task=" + task + "]";
|
||||
}
|
||||
}
|
||||
try {
|
||||
Future<?> f = executor.submit(sr);
|
||||
sr.set(f);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
sr.dispose();
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
return sr;
|
||||
}
|
||||
|
||||
void delete(CancelFuture r) {
|
||||
synchronized (this) {
|
||||
if (!terminated) {
|
||||
tasks.remove(r);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedule(Runnable task, long delay, TimeUnit unit) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
TimedScheduledRunnable sr = new TimedScheduledRunnable(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
tasks.add(sr);
|
||||
}
|
||||
|
||||
try {
|
||||
Future<?> f = executor.schedule(sr, delay, unit);
|
||||
sr.set(f);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
sr.dispose();
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
return sr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Cancellation schedulePeriodically(Runnable task, long initialDelay, long period, TimeUnit unit) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
TimedPeriodicScheduledRunnable sr = new TimedPeriodicScheduledRunnable(task, this);
|
||||
|
||||
synchronized (this) {
|
||||
if (terminated) {
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
tasks.add(sr);
|
||||
}
|
||||
|
||||
try {
|
||||
Future<?> f = executor.scheduleAtFixedRate(sr, initialDelay, period, unit);
|
||||
sr.set(f);
|
||||
} catch (RejectedExecutionException ex) {
|
||||
sr.dispose();
|
||||
return REJECTED;
|
||||
}
|
||||
|
||||
return sr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
if (terminated) {
|
||||
return;
|
||||
}
|
||||
terminated = true;
|
||||
|
||||
OpenHashSet<CancelFuture> set;
|
||||
|
||||
synchronized (this) {
|
||||
set = tasks;
|
||||
if (set == null) {
|
||||
return;
|
||||
}
|
||||
tasks = null;
|
||||
}
|
||||
|
||||
if (!set.isEmpty()) {
|
||||
Object[] keys = set.keys();
|
||||
for (Object c : keys) {
|
||||
if (c != null) {
|
||||
((CancelFuture) c).cancelFuture();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
interface CancelFuture {
|
||||
void cancelFuture();
|
||||
}
|
||||
|
||||
static final class TimedScheduledRunnable
|
||||
extends AtomicReference<Future<?>> implements Runnable, Cancellation, CancelFuture {
|
||||
/** */
|
||||
private static final long serialVersionUID = 2284024836904862408L;
|
||||
|
||||
final Runnable task;
|
||||
|
||||
final SingleTimedSchedulerWorker parent;
|
||||
|
||||
volatile Thread current;
|
||||
static final AtomicReferenceFieldUpdater<TimedScheduledRunnable, Thread> CURRENT =
|
||||
AtomicReferenceFieldUpdater.newUpdater(TimedScheduledRunnable.class, Thread.class, "current");
|
||||
|
||||
static final Runnable EMPTY = new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
|
||||
}
|
||||
};
|
||||
|
||||
static final Future<?> CANCELLED_FUTURE = new FutureTask<>(EMPTY, null);
|
||||
|
||||
static final Future<?> FINISHED = new FutureTask<>(EMPTY, null);
|
||||
|
||||
public TimedScheduledRunnable(Runnable task, SingleTimedSchedulerWorker parent) {
|
||||
this.task = task;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
CURRENT.lazySet(this, Thread.currentThread());
|
||||
try {
|
||||
try {
|
||||
task.run();
|
||||
} catch (Throwable e) {
|
||||
Operators.onErrorDropped(e);
|
||||
}
|
||||
} finally {
|
||||
for (; ; ) {
|
||||
Future<?> a = get();
|
||||
if (a == CANCELLED_FUTURE) {
|
||||
break;
|
||||
}
|
||||
if (compareAndSet(a, FINISHED)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
parent.delete(this);
|
||||
break;
|
||||
}
|
||||
}
|
||||
CURRENT.lazySet(this, null);
|
||||
}
|
||||
}
|
||||
|
||||
void doCancel(Future<?> a) {
|
||||
a.cancel(Thread.currentThread() != current);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void cancelFuture() {
|
||||
for (; ; ) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(a, CANCELLED_FUTURE)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
for (; ; ) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(a, CANCELLED_FUTURE)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
parent.delete(this);
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
for (; ; ) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (a == CANCELLED_FUTURE) {
|
||||
doCancel(a);
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(null, f)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "TimedScheduledRunnable[cancelled=" + (get() == CANCELLED_FUTURE) +
|
||||
", task=" + task +
|
||||
"]";
|
||||
}
|
||||
}
|
||||
|
||||
static final class TimedPeriodicScheduledRunnable
|
||||
extends AtomicReference<Future<?>>
|
||||
implements Runnable, Cancellation, CancelFuture {
|
||||
/** */
|
||||
private static final long serialVersionUID = 2284024836904862408L;
|
||||
|
||||
final Runnable task;
|
||||
|
||||
final SingleTimedSchedulerWorker parent;
|
||||
|
||||
volatile Thread current;
|
||||
static final AtomicReferenceFieldUpdater<TimedPeriodicScheduledRunnable, Thread> CURRENT =
|
||||
AtomicReferenceFieldUpdater.newUpdater(TimedPeriodicScheduledRunnable.class, Thread.class, "current");
|
||||
|
||||
static final Runnable EMPTY = new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
|
||||
}
|
||||
};
|
||||
|
||||
static final Future<?> CANCELLED_FUTURE = new FutureTask<>(EMPTY, null);
|
||||
|
||||
static final Future<?> FINISHED = new FutureTask<>(EMPTY, null);
|
||||
|
||||
public TimedPeriodicScheduledRunnable(Runnable task, SingleTimedSchedulerWorker parent) {
|
||||
this.task = task;
|
||||
this.parent = parent;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
CURRENT.lazySet(this, Thread.currentThread());
|
||||
try {
|
||||
try {
|
||||
task.run();
|
||||
} catch (Throwable ex) {
|
||||
Operators.onErrorDropped(ex);
|
||||
for (; ; ) {
|
||||
Future<?> a = get();
|
||||
if (a == CANCELLED_FUTURE) {
|
||||
break;
|
||||
}
|
||||
if (compareAndSet(a, FINISHED)) {
|
||||
parent.delete(this);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
CURRENT.lazySet(this, null);
|
||||
}
|
||||
}
|
||||
|
||||
void doCancel(Future<?> a) {
|
||||
a.cancel(false);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void cancelFuture() {
|
||||
for (; ; ) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(a, CANCELLED_FUTURE)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void dispose() {
|
||||
for (; ; ) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(a, CANCELLED_FUTURE)) {
|
||||
if (a != null) {
|
||||
doCancel(a);
|
||||
}
|
||||
parent.delete(this);
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
void setFuture(Future<?> f) {
|
||||
for (; ; ) {
|
||||
Future<?> a = get();
|
||||
if (a == FINISHED) {
|
||||
return;
|
||||
}
|
||||
if (a == CANCELLED_FUTURE) {
|
||||
doCancel(a);
|
||||
return;
|
||||
}
|
||||
if (compareAndSet(null, f)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "TimedPeriodicScheduledRunnable[cancelled=" + get() + ", task=" + task + "]";
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user