This is an automated email from the ASF dual-hosted git repository.

clebertsuconic pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/activemq-artemis.git


The following commit(s) were added to refs/heads/main by this push:
     new 9acd036dcc ARTEMIS-3928 Fixing yield and shutdownNow
9acd036dcc is described below

commit 9acd036dcc9e76b19797817cc7bdc30a29706a48
Author: Clebert Suconic <[email protected]>
AuthorDate: Fri Aug 19 15:55:08 2022 -0400

    ARTEMIS-3928 Fixing yield and shutdownNow
    
    
org.apache.activemq.artemis.tests.integration.jms.client.ReceiveNoWaitTest.testReceiveNoWait
 was failing because before this change
---
 .../artemis/utils/actors/ProcessorBase.java        | 38 +++++++++++-----------
 1 file changed, 19 insertions(+), 19 deletions(-)

diff --git 
a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/actors/ProcessorBase.java
 
b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/actors/ProcessorBase.java
index b631f52906..9a27de7fac 100644
--- 
a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/actors/ProcessorBase.java
+++ 
b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/actors/ProcessorBase.java
@@ -40,19 +40,17 @@ public abstract class ProcessorBase<T> extends HandlerBase {
     * Using a method reference instead of an inner classes allows the caller 
to reduce the pointer chasing
     * when accessing ProcessorBase.this fields/methods.
     */
-   private final Runnable mainTask = this::executePendingTasks;
+   private final Runnable task = this::executePendingTasks;
 
    // used by stateUpdater
    @SuppressWarnings("unused")
    private volatile int state = STATE_NOT_RUNNING;
-
-   private enum request {
-      keepRunning,
-      shutdown,
-      yield
-   }
-
-   private volatile request loopRequest = request.keepRunning;
+   // Request of forced shutdown
+   private volatile boolean requestedForcedShutdown = false;
+   // Request of educated shutdown:
+   private volatile boolean requestedShutdown = false;
+   // Request to yield to another thread
+   private volatile boolean yielded = false;
 
    private static final AtomicIntegerFieldUpdater<ProcessorBase> stateUpdater 
= AtomicIntegerFieldUpdater.newUpdater(ProcessorBase.class, "state");
 
@@ -65,7 +63,7 @@ public abstract class ProcessorBase<T> extends HandlerBase {
                T task;
                //while the queue is not empty we process in order:
                //if requestedForcedShutdown==true than no new tasks will be 
drained from the tasks q.
-               while (loopRequest == request.keepRunning && (task = 
tasks.poll()) != null) {
+               while (!yielded && !requestedForcedShutdown && (task = 
tasks.poll()) != null) {
                   doTask(task);
                }
             } finally {
@@ -83,11 +81,11 @@ public abstract class ProcessorBase<T> extends HandlerBase {
          //but poll() has returned null, so a submitting thread will believe 
that it does not need re-execute.
          //this check fixes the issue
       }
-      while (!tasks.isEmpty() && loopRequest == request.keepRunning);
+      while (!tasks.isEmpty() && !requestedShutdown && !yielded);
 
-      if (loopRequest == request.yield) {
-         loopRequest = request.keepRunning;
-         delegate.execute(mainTask);
+      if (yielded) {
+         yielded = false;
+         delegate.execute(task);
       }
    }
 
@@ -99,7 +97,7 @@ public abstract class ProcessorBase<T> extends HandlerBase {
    }
 
    public void shutdown(long timeout, TimeUnit unit) {
-      loopRequest = request.shutdown;
+      requestedShutdown = true;
 
       if (!inHandler()) {
          // if it's in handler.. we just return
@@ -108,13 +106,15 @@ public abstract class ProcessorBase<T> extends 
HandlerBase {
    }
 
    public void yield() {
-      this.loopRequest = request.yield;
+      this.yielded = true;
    }
 
    /** It will shutdown the executor however it will not wait for finishing 
tasks*/
    public int shutdownNow(Consumer<? super T> onPendingItem, int timeout, 
TimeUnit unit) {
       //alert anyone that has been requested (at least) an immediate shutdown
-      loopRequest = request.shutdown;
+      requestedForcedShutdown = true;
+      requestedShutdown = true;
+      yielded = false;
 
       if (!inHandler()) {
          // We don't have an option where we could do an immediate timeout
@@ -174,7 +174,7 @@ public abstract class ProcessorBase<T> extends HandlerBase {
    }
 
    protected void task(T command) {
-      if (loopRequest == request.shutdown) {
+      if (requestedShutdown) {
          logAddOnShutdown();
          return;
       }
@@ -195,7 +195,7 @@ public abstract class ProcessorBase<T> extends HandlerBase {
    private void onAddedTaskIfNotRunning(int state) {
       if (state == STATE_NOT_RUNNING) {
          //startPoller could be deleted but is maintained because is inherited
-         delegate.execute(mainTask);
+         delegate.execute(task);
       }
    }
 

Reply via email to