From 512a5d9190d8f8686374d9f6428e922b16692c25 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Thu, 25 Aug 2016 15:27:28 -0400 Subject: [PATCH] Use spring-cloud-build-tools for checkstyle validation --- pom.xml | 6 +- .../core/scheduler/ElasticScheduler.java | 555 ++++++------ .../core/scheduler/ParallelScheduler.java | 550 ++++++------ .../core/scheduler/SingleScheduler.java | 496 +++++------ .../core/scheduler/SingleTimedScheduler.java | 813 +++++++++--------- 5 files changed, 1213 insertions(+), 1207 deletions(-) diff --git a/pom.xml b/pom.xml index 0d2d0f7c5..0e1975adf 100644 --- a/pom.xml +++ b/pom.xml @@ -150,8 +150,8 @@ org.springframework.cloud - spring-cloud-stream-tools - 1.1.0.BUILD-SNAPSHOT + spring-cloud-build-tools + 1.2.0.RELEASE @@ -160,8 +160,6 @@ validate checkstyle.xml - checkstyle-header.txt - checkstyle-suppressions.xml UTF-8 true true diff --git a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/ElasticScheduler.java b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/ElasticScheduler.java index 8101de45b..2403dc364 100644 --- a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/ElasticScheduler.java +++ b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/ElasticScheduler.java @@ -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 cache; + final ThreadFactory factory; - final Queue 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 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 list = new ArrayList<>(cache); - for (ExecutorServiceExpiry e : list) { - if (e.expireMillis < now) { - if (cache.remove(e)) { - e.executor.shutdownNow(); - } - } - } - } + final Queue 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 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 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> - 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 CANCELLED = new FutureTask<>(() -> { }, null); + cache.clear(); - static final FutureTask 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 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 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 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> + implements Runnable, Cancellation { + /** */ + private static final long serialVersionUID = 6799295393954430738L; + + final Runnable run; + + final CachedWorker parent; + + volatile boolean cancelled; + + static final FutureTask CANCELLED = new FutureTask<>(() -> { + }, null); + + static final FutureTask 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); + } + } + } + } + } } diff --git a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/ParallelScheduler.java b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/ParallelScheduler.java index a062dfcff..ad72beb21 100644 --- a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/ParallelScheduler.java +++ b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/ParallelScheduler.java @@ -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 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 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 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 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 FUTURE = - AtomicReferenceFieldUpdater.newUpdater(ParallelWorkerTask.class, Future.class, "future"); - - static final Future FINISHED = CompletableFuture.completedFuture(null); - static final Future 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 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 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 FUTURE = + AtomicReferenceFieldUpdater.newUpdater(ParallelWorkerTask.class, Future.class, "future"); + + static final Future FINISHED = CompletableFuture.completedFuture(null); + static final Future 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); + } + } + } + } + } } diff --git a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/SingleScheduler.java b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/SingleScheduler.java index 55bfdee05..aa3b77fc6 100644 --- a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/SingleScheduler.java +++ b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/SingleScheduler.java @@ -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 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 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 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 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 FUTURE = - AtomicReferenceFieldUpdater.newUpdater(SingleWorkerTask.class, Future.class, "future"); - - static final Future FINISHED = CompletableFuture.completedFuture(null); - static final Future 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 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 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 FUTURE = + AtomicReferenceFieldUpdater.newUpdater(SingleWorkerTask.class, Future.class, "future"); + + static final Future FINISHED = CompletableFuture.completedFuture(null); + static final Future 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); + } + } + } + } + } } diff --git a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/SingleTimedScheduler.java b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/SingleTimedScheduler.java index 1aecc49bf..33ffdc8ce 100644 --- a/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/SingleTimedScheduler.java +++ b/spring-cloud-stream-reactive/src/main/java/org/springframework/cloud/stream/reactive/reactor/core/scheduler/SingleTimedScheduler.java @@ -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 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 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> implements Runnable, Cancellation, CancelFuture { - /** */ - private static final long serialVersionUID = 2284024836904862408L; - - final Runnable task; - - final SingleTimedSchedulerWorker parent; - - volatile Thread current; - static final AtomicReferenceFieldUpdater 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 tasks; - static final class TimedPeriodicScheduledRunnable - extends AtomicReference> - implements Runnable, Cancellation, CancelFuture { - /** */ - private static final long serialVersionUID = 2284024836904862408L; - - final Runnable task; - - final SingleTimedSchedulerWorker parent; - - volatile Thread current; - static final AtomicReferenceFieldUpdater 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 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> implements Runnable, Cancellation, CancelFuture { + /** */ + private static final long serialVersionUID = 2284024836904862408L; + + final Runnable task; + + final SingleTimedSchedulerWorker parent; + + volatile Thread current; + static final AtomicReferenceFieldUpdater 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> + implements Runnable, Cancellation, CancelFuture { + /** */ + private static final long serialVersionUID = 2284024836904862408L; + + final Runnable task; + + final SingleTimedSchedulerWorker parent; + + volatile Thread current; + static final AtomicReferenceFieldUpdater 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 + "]"; + } + } }