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

vavrtom pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/qpid-broker-j.git


The following commit(s) were added to refs/heads/main by this push:
     new e901f67448 QPID-8752: [Broker-J] Connection- and session-level Subject 
caching (#429)
e901f67448 is described below

commit e901f674483918d767ae9b631857ea08c27c1f86
Author: Daniil Kirilyuk <[email protected]>
AuthorDate: Wed Sep 2 15:20:43 2026 +0200

    QPID-8752: [Broker-J] Connection- and session-level Subject caching (#429)
---
 broker-core/pom.xml                                |  12 +-
 .../server/security/SubjectExecutionContext.java   |  91 ++++++++++++-
 .../qpid/server/session/AbstractAMQPSession.java   |   8 ++
 .../server/transport/AbstractAMQPConnection.java   |  15 ++-
 .../server/security/SubjectExecutionContext.java   |  28 +++-
 .../security/SubjectExecutionContextTest.java      | 142 ++++++++++++++++++++-
 .../server/protocol/v0_10/ServerAssembler.java     |   3 +-
 .../qpid/server/protocol/v0_10/ServerSession.java  |   7 +-
 .../qpid/server/protocol/v0_8/AMQChannel.java      |   5 +-
 .../protocol/v0_8/AMQPConnection_0_8Impl.java      |   5 +-
 .../qpid/server/protocol/v0_8/BrokerDecoder.java   |   2 +-
 .../protocol/v1_0/AMQPConnection_1_0Impl.java      |  12 +-
 .../qpid/server/protocol/v1_0/Session_1_0.java     |   2 +-
 13 files changed, 297 insertions(+), 35 deletions(-)

diff --git a/broker-core/pom.xml b/broker-core/pom.xml
index f4f31aa051..f8248a6cd6 100644
--- a/broker-core/pom.xml
+++ b/broker-core/pom.xml
@@ -241,7 +241,7 @@
     <profile>
       <id>jdk17on</id>
       <activation>
-        <jdk>(,23]</jdk>
+        <jdk>(,23)</jdk>
       </activation>
       <build>
         <plugins>
@@ -261,9 +261,9 @@
       </build>
     </profile>
     <profile>
-      <id>jdk24on</id>
+      <id>jdk23on</id>
       <activation>
-        <jdk>[24,)</jdk>
+        <jdk>[23,)</jdk>
       </activation>
       <build>
         <plugins>
@@ -272,14 +272,14 @@
             <artifactId>maven-compiler-plugin</artifactId>
             <executions>
               <execution>
-                <id>java24-compile</id>
+                <id>java23-compile</id>
                 <phase>compile</phase>
                 <goals>
                   <goal>compile</goal>
                 </goals>
                 <configuration>
-                  <release>24</release>
-                  
<compileSourceRoots>${project.basedir}/src/main/java24</compileSourceRoots>
+                  <release>23</release>
+                  
<compileSourceRoots>${project.basedir}/src/main/java23</compileSourceRoots>
                   <multiReleaseOutput>true</multiReleaseOutput>
                 </configuration>
               </execution>
diff --git 
a/broker-core/src/main/java/org/apache/qpid/server/security/SubjectExecutionContext.java
 
b/broker-core/src/main/java/org/apache/qpid/server/security/SubjectExecutionContext.java
index 1e2c30d318..3873771fe7 100644
--- 
a/broker-core/src/main/java/org/apache/qpid/server/security/SubjectExecutionContext.java
+++ 
b/broker-core/src/main/java/org/apache/qpid/server/security/SubjectExecutionContext.java
@@ -33,15 +33,25 @@ import java.util.function.Function;
 import javax.security.auth.Subject;
 
 /**
- * Java 17-23 implementation backed by {@link 
Subject#getSubject(AccessControlContext)} ()} and {@link Subject#doAs(Subject, 
PrivilegedAction)}}.
+ * Java 17-22 implementation backed by {@link 
Subject#getSubject(AccessControlContext)} and the deprecated access
+ * control APIs.
  * <br>
- * Provides the same API surface as the Java 24+ implementation, but relies on 
the deprecated Security manager APIs.
+ * Provides the same API surface as the Java 23+ implementation. A reusable 
instance captures the caller's access
+ * control context and combines it with the supplied subject once, avoiding 
the allocation performed by
+ * {@link Subject#doAs(Subject, PrivilegedAction)} on every invocation.
  */
 public final class SubjectExecutionContext
 {
-    private SubjectExecutionContext()
+    private final AccessControlContext _accessControlContext;
+
+    private SubjectExecutionContext(final Subject subject)
     {
-        // utility class has private constructor
+        _accessControlContext = Subject.doAs(subject, 
(PrivilegedAction<AccessControlContext>) AccessController::getContext);
+    }
+
+    public static SubjectExecutionContext create(final Subject subject)
+    {
+        return new SubjectExecutionContext(subject);
     }
 
     public static Subject currentSubject()
@@ -127,6 +137,79 @@ public final class SubjectExecutionContext
         });
     }
 
+    public <T> T call(final Callable<T> action) throws Exception
+    {
+        try
+        {
+            return 
AccessController.doPrivileged((PrivilegedExceptionAction<T>) action::call, 
_accessControlContext);
+        }
+        catch (PrivilegedActionException pae)
+        {
+            final Throwable cause = pae.getCause();
+            if (cause == null)
+            {
+                throw pae;
+            }
+            if (cause instanceof Error err)
+            {
+                err.addSuppressed(pae);
+                throw err;
+            }
+            if (cause instanceof RuntimeException re)
+            {
+                re.addSuppressed(pae);
+                throw re;
+            }
+            if (cause instanceof Exception ex)
+            {
+                ex.addSuppressed(pae);
+                throw ex;
+            }
+            throw pae;
+        }
+    }
+
+    public <T> T callUnchecked(final Callable<T> action)
+    {
+        try
+        {
+            return 
AccessController.doPrivileged((PrivilegedExceptionAction<T>) action::call, 
_accessControlContext);
+        }
+        catch (PrivilegedActionException pae)
+        {
+            final Throwable cause = pae.getCause();
+            if (cause == null)
+            {
+                throw new CompletionException(pae);
+            }
+            if (cause instanceof Error err)
+            {
+                err.addSuppressed(pae);
+                throw err;
+            }
+            if (cause instanceof RuntimeException re)
+            {
+                re.addSuppressed(pae);
+                throw re;
+            }
+            if (cause instanceof Exception ex)
+            {
+                ex.addSuppressed(pae);
+                throw new SubjectActionException(ex);
+            }
+            throw new CompletionException(pae);
+        }
+    }
+
+    public void run(final Runnable action)
+    {
+        AccessController.doPrivileged((PrivilegedAction<Object>) () ->
+        {
+            action.run();
+            return null;
+        }, _accessControlContext);
+    }
+
     public static Throwable unwrapSubjectActionException(final Throwable 
throwable)
     {
         return ((throwable instanceof SubjectActionException sae && 
sae.getCause() != null)
diff --git 
a/broker-core/src/main/java/org/apache/qpid/server/session/AbstractAMQPSession.java
 
b/broker-core/src/main/java/org/apache/qpid/server/session/AbstractAMQPSession.java
index c529b913c4..5162f21365 100644
--- 
a/broker-core/src/main/java/org/apache/qpid/server/session/AbstractAMQPSession.java
+++ 
b/broker-core/src/main/java/org/apache/qpid/server/session/AbstractAMQPSession.java
@@ -63,6 +63,7 @@ import org.apache.qpid.server.model.Session;
 import org.apache.qpid.server.model.State;
 import org.apache.qpid.server.protocol.PublishAuthorisationCache;
 import org.apache.qpid.server.security.SecurityToken;
+import org.apache.qpid.server.security.SubjectExecutionContext;
 import org.apache.qpid.server.transport.AMQPConnection;
 import org.apache.qpid.server.transport.network.Ticker;
 import org.apache.qpid.server.util.Action;
@@ -77,6 +78,7 @@ public abstract class AbstractAMQPSession<S extends 
AbstractAMQPSession<S, X>,
     private final Action _deleteModelTask;
     private final AMQPConnection<?> _connection;
     private final int _sessionId;
+    private final SubjectExecutionContext _subjectExecutionContext;
 
     protected final Subject _subject;
     protected final SecurityToken _token;
@@ -132,6 +134,7 @@ public abstract class AbstractAMQPSession<S extends 
AbstractAMQPSession<S, X>,
             final Broker<?> broker = (Broker<?>) _connection.getBroker();
             _token = broker.newToken(_subject);
         }
+        _subjectExecutionContext = SubjectExecutionContext.create(_subject);
 
         final long authCacheTimeout = _connection.getContextValue(Long.class, 
Session.PRODUCER_AUTH_CACHE_TIMEOUT);
         final int authCacheSize = _connection.getContextValue(Integer.class, 
Session.PRODUCER_AUTH_CACHE_SIZE);
@@ -169,6 +172,11 @@ public abstract class AbstractAMQPSession<S extends 
AbstractAMQPSession<S, X>,
         return _connection;
     }
 
+    public final SubjectExecutionContext getSubjectExecutionContext()
+    {
+        return _subjectExecutionContext;
+    }
+
     @Override
     public boolean isProducerFlowBlocked()
     {
diff --git 
a/broker-core/src/main/java/org/apache/qpid/server/transport/AbstractAMQPConnection.java
 
b/broker-core/src/main/java/org/apache/qpid/server/transport/AbstractAMQPConnection.java
index 1592e93dab..88ab2f8c38 100644
--- 
a/broker-core/src/main/java/org/apache/qpid/server/transport/AbstractAMQPConnection.java
+++ 
b/broker-core/src/main/java/org/apache/qpid/server/transport/AbstractAMQPConnection.java
@@ -104,6 +104,7 @@ public abstract class AbstractAMQPConnection<C extends 
AbstractAMQPConnection<C,
     private final List<Action<? super C>> _connectionCloseTaskList = new 
CopyOnWriteArrayList<>();
 
     private final LogSubject _logSubject;
+    private volatile SubjectExecutionContext _subjectExecutionContext;
     private volatile ContextProvider _contextProvider;
     private volatile EventLoggerProvider _eventLoggerProvider;
     private String _clientProduct;
@@ -167,6 +168,7 @@ public abstract class AbstractAMQPConnection<C extends 
AbstractAMQPConnection<C,
         _aggregateTicker = aggregateTicker;
         _subject = new Subject();
         _subject.getPrincipals().add(new ConnectionPrincipal(this));
+        updateSubjectExecutionContext();
 
         _transportClosedFuture.thenRunAsync(() ->
         {
@@ -535,7 +537,7 @@ public abstract class AbstractAMQPConnection<C extends 
AbstractAMQPConnection<C,
     @Override
     public final void received(final QpidByteBuffer buf)
     {
-        SubjectExecutionContext.withSubject(_subject, () ->
+        _subjectExecutionContext.run(() ->
         {
             updateLastReadTime();
             try
@@ -562,9 +564,9 @@ public abstract class AbstractAMQPConnection<C extends 
AbstractAMQPConnection<C,
 
     protected abstract boolean isOpeningInProgress();
 
-    protected <T> T runAsSubject(Supplier<T> action)
+    protected <T> T runAsSubject(final Supplier<T> action)
     {
-        return SubjectExecutionContext.withSubjectUnchecked(_subject, 
action::get);
+        return _subjectExecutionContext.callUnchecked(action::get);
     }
 
     private boolean runningAsSubject()
@@ -856,6 +858,7 @@ public abstract class AbstractAMQPConnection<C extends 
AbstractAMQPConnection<C,
 
         _subject.getPrincipals().add(addressSpace.getPrincipal());
 
+        updateSubjectExecutionContext();
         logConnectionOpen();
     }
 
@@ -893,6 +896,12 @@ public abstract class AbstractAMQPConnection<C extends 
AbstractAMQPConnection<C,
         _subject.getPrincipals().addAll(subject.getPrincipals());
         
_subject.getPrivateCredentials().addAll(subject.getPrivateCredentials());
         _subject.getPublicCredentials().addAll(subject.getPublicCredentials());
+        updateSubjectExecutionContext();
+    }
+
+    private void updateSubjectExecutionContext()
+    {
+        _subjectExecutionContext = SubjectExecutionContext.create(_subject);
     }
 
     @Override
diff --git 
a/broker-core/src/main/java24/org/apache/qpid/server/security/SubjectExecutionContext.java
 
b/broker-core/src/main/java23/org/apache/qpid/server/security/SubjectExecutionContext.java
similarity index 87%
rename from 
broker-core/src/main/java24/org/apache/qpid/server/security/SubjectExecutionContext.java
rename to 
broker-core/src/main/java23/org/apache/qpid/server/security/SubjectExecutionContext.java
index c0dd6e0749..666337edb3 100644
--- 
a/broker-core/src/main/java24/org/apache/qpid/server/security/SubjectExecutionContext.java
+++ 
b/broker-core/src/main/java23/org/apache/qpid/server/security/SubjectExecutionContext.java
@@ -28,16 +28,23 @@ import java.util.function.Function;
 import javax.security.auth.Subject;
 
 /**
- * Java 24+ implementation backed by {@link Subject#current()} and {@link 
Subject#callAs(Subject, Callable)}.
+ * Java 23+ implementation backed by {@link Subject#current()} and {@link 
Subject#callAs(Subject, Callable)}.
  * <br>
  * Provides the same API surface as the Java 17 implementation, but relies on 
the JDK-managed current Subject
  * rather than removed Subject.getSubject().
  */
 public final class SubjectExecutionContext
 {
-    private SubjectExecutionContext()
+    private final Subject _subject;
+
+    private SubjectExecutionContext(final Subject subject)
+    {
+        _subject = subject;
+    }
+
+    public static SubjectExecutionContext create(final Subject subject)
     {
-        // utility class has private constructor
+        return new SubjectExecutionContext(subject);
     }
 
     public static Subject currentSubject()
@@ -145,6 +152,21 @@ public final class SubjectExecutionContext
         }
     }
 
+    public <T> T call(final Callable<T> action) throws Exception
+    {
+        return withSubject(_subject, action);
+    }
+
+    public <T> T callUnchecked(final Callable<T> action)
+    {
+        return withSubjectUnchecked(_subject, action);
+    }
+
+    public void run(final Runnable action)
+    {
+        withSubject(_subject, action);
+    }
+
     public static Throwable unwrapSubjectActionException(final Throwable 
throwable)
     {
         return ((throwable instanceof SubjectActionException sae && 
sae.getCause() != null)
diff --git 
a/broker-core/src/test/java/org/apache/qpid/server/security/SubjectExecutionContextTest.java
 
b/broker-core/src/test/java/org/apache/qpid/server/security/SubjectExecutionContextTest.java
index eb2f67de03..196fe5f75a 100644
--- 
a/broker-core/src/test/java/org/apache/qpid/server/security/SubjectExecutionContextTest.java
+++ 
b/broker-core/src/test/java/org/apache/qpid/server/security/SubjectExecutionContextTest.java
@@ -22,25 +22,36 @@
 package org.apache.qpid.server.security;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
+import java.security.Permission;
+import java.security.Policy;
+import java.security.ProtectionDomain;
+import java.security.SecurityPermission;
+import java.util.ArrayList;
+import java.util.List;
 import java.util.concurrent.Callable;
 import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
 
 import javax.security.auth.Subject;
 
 import org.junit.jupiter.api.Test;
-
-import org.apache.qpid.test.utils.UnitTestBase;
 import org.junit.jupiter.api.condition.EnabledForJreRange;
 import org.junit.jupiter.api.condition.JRE;
 
+import org.apache.qpid.test.utils.UnitTestBase;
+
 class SubjectExecutionContextTest extends UnitTestBase
 {
     private static final String SUBJECT_MUST_BE_NULL_OUTSIDE_CONTEXT = 
"Subject must be null outside context";
@@ -88,6 +99,133 @@ class SubjectExecutionContextTest extends UnitTestBase
         assertNull(SubjectExecutionContext.currentSubject(), 
SUBJECT_MUST_BE_RESTORED_AFTER_CONTEXT);
     }
 
+    @Test
+    void reusableContextRestoresPreviousSubject()
+    {
+        final Subject subjectA = new Subject();
+        final Subject subjectB = new Subject();
+        final SubjectExecutionContext contextA = 
SubjectExecutionContext.create(subjectA);
+        final SubjectExecutionContext contextB = 
SubjectExecutionContext.create(subjectB);
+
+        assertNull(SubjectExecutionContext.currentSubject());
+
+        contextA.run(() ->
+        {
+            assertSame(subjectA, SubjectExecutionContext.currentSubject());
+            contextB.run(() -> assertSame(subjectB, 
SubjectExecutionContext.currentSubject()));
+            assertSame(subjectA, SubjectExecutionContext.currentSubject());
+        });
+
+        assertNull(SubjectExecutionContext.currentSubject());
+    }
+
+    @Test
+    @EnabledForJreRange(max = JRE.JAVA_17)
+    void 
reusableContextCreationDoesNotRequireCreateAccessControlContextPermission()
+    {
+        final Policy originalPolicy = Policy.getPolicy();
+        final SecurityManager originalSecurityManager = 
System.getSecurityManager();
+        try
+        {
+            Policy.setPolicy(new Policy()
+            {
+                @Override
+                public boolean implies(final ProtectionDomain domain, final 
Permission permission)
+                {
+                    return !(permission instanceof SecurityPermission &&
+                            
"createAccessControlContext".equals(permission.getName()));
+                }
+            });
+            System.setSecurityManager(new SecurityManager());
+
+            final Subject subject = new Subject();
+            final SubjectExecutionContext context = 
SubjectExecutionContext.create(subject);
+
+            context.run(() -> assertSame(subject, 
SubjectExecutionContext.currentSubject()));
+        }
+        finally
+        {
+            System.setSecurityManager(originalSecurityManager);
+            Policy.setPolicy(originalPolicy);
+        }
+    }
+
+    @Test
+    void reusableContextCallRestoresAfterException()
+    {
+        final Subject subject = new Subject();
+        final SubjectExecutionContext context = 
SubjectExecutionContext.create(subject);
+        final Exception expected = new Exception();
+
+        final Exception thrown = assertThrows(Exception.class, () -> 
context.call(() ->
+        {
+            assertSame(subject, SubjectExecutionContext.currentSubject());
+            throw expected;
+        }));
+
+        assertSame(expected, thrown);
+        assertNull(SubjectExecutionContext.currentSubject());
+    }
+
+    @Test
+    void reusableContextCallUncheckedWrapsCheckedException()
+    {
+        final SubjectExecutionContext context = 
SubjectExecutionContext.create(new Subject());
+        final Exception expected = new Exception();
+
+        final SubjectActionException thrown = 
assertThrows(SubjectActionException.class, () ->
+                context.callUnchecked(() ->
+                {
+                    throw expected;
+                }));
+
+        assertSame(expected, thrown.getCause());
+        assertNull(SubjectExecutionContext.currentSubject());
+    }
+
+    @Test
+    void reusableContextSupportsConcurrentUse() throws Exception
+    {
+        final long timeout = 10L;
+        final int concurrencyLevel = 4;
+        final CyclicBarrier barrier = new CyclicBarrier(concurrencyLevel);
+        final Subject subject = new Subject();
+        final SubjectExecutionContext context = 
SubjectExecutionContext.create(subject);
+        final ExecutorService executor = 
Executors.newFixedThreadPool(concurrencyLevel);
+        final List<Callable<Subject>> actions = new ArrayList<>();
+        for (int i = 0; i < concurrencyLevel; i++)
+        {
+            actions.add(() ->
+            {
+                final Subject observedSubject = context.call(() ->
+                {
+                    barrier.await(timeout, TimeUnit.SECONDS);
+                    return SubjectExecutionContext.currentSubject();
+                });
+
+                assertNull(SubjectExecutionContext.currentSubject(), "Subject 
must be restored after context");
+
+                return observedSubject;
+            });
+        }
+
+        try
+        {
+            final List<Future<Subject>> results = executor.invokeAll(actions, 
timeout, TimeUnit.SECONDS);
+            for (final Future<Subject> result : results)
+            {
+                assertFalse(result.isCancelled(), "Concurrent subject action 
timed out");
+                assertSame(subject, result.get());
+            }
+        }
+        finally
+        {
+            executor.shutdownNow();
+        }
+
+        assertNull(SubjectExecutionContext.currentSubject(), "Subject must be 
restored after context");
+    }
+
     @Test
     void withSubjectRestoresAfterException()
     {
diff --git 
a/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerAssembler.java
 
b/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerAssembler.java
index d98a10271b..93457ddf3e 100644
--- 
a/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerAssembler.java
+++ 
b/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerAssembler.java
@@ -40,7 +40,6 @@ import 
org.apache.qpid.server.protocol.v0_10.transport.ProtocolError;
 import org.apache.qpid.server.protocol.v0_10.transport.ProtocolEvent;
 import org.apache.qpid.server.protocol.v0_10.transport.ProtocolHeader;
 import org.apache.qpid.server.protocol.v0_10.transport.Struct;
-import org.apache.qpid.server.security.SubjectExecutionContext;
 import org.apache.qpid.server.util.PeekingIterator;
 import org.apache.qpid.server.util.PeekingIteratorImpl;
 
@@ -84,7 +83,7 @@ public class ServerAssembler
                     final ServerSession channel = 
_connection.getSession(frameChannel);
                     if (channel != null)
                     {
-                        
SubjectExecutionContext.withSubject(channel.getSubject(), () ->
+                        channel.getSubjectExecutionContext().run(() ->
                         {
                             ServerFrame channelFrame = frame;
                             boolean nextIsSameChannel;
diff --git 
a/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
 
b/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
index b3444a5b8f..1c4827f741 100644
--- 
a/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
+++ 
b/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
@@ -894,6 +894,11 @@ public class ServerSession extends SessionInvoker
         return _modelObject.getSubject();
     }
 
+    public SubjectExecutionContext getSubjectExecutionContext()
+    {
+        return _modelObject.getSubjectExecutionContext();
+    }
+
     protected void setState(final State state)
     {
         if(runningAsSubject())
@@ -928,7 +933,7 @@ public class ServerSession extends SessionInvoker
 
     private <T> T runAsSubject(final Supplier<T> action)
     {
-        return 
SubjectExecutionContext.withSubjectUnchecked(getAuthorizedSubject(), 
action::get);
+        return getSubjectExecutionContext().callUnchecked(action::get);
     }
 
     private void invokeBlock()
diff --git 
a/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
 
b/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
index 2f979f50dd..0b2fba51e9 100644
--- 
a/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
+++ 
b/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java
@@ -98,7 +98,6 @@ import org.apache.qpid.server.txn.LocalTransaction;
 import org.apache.qpid.server.txn.ServerTransaction;
 import org.apache.qpid.server.util.Action;
 import org.apache.qpid.server.security.AccessDeniedException;
-import org.apache.qpid.server.security.SubjectExecutionContext;
 import 
org.apache.qpid.server.virtualhost.MessageDestinationIsAlternateException;
 import org.apache.qpid.server.virtualhost.RequiredExchangeException;
 import org.apache.qpid.server.virtualhost.ReservedExchangeNameException;
@@ -219,7 +218,7 @@ public class AMQChannel extends 
AbstractAMQPSession<AMQChannel, ConsumerTarget_0
 
         _clientDeliveryMethod = connection.createDeliveryMethod(_channelId);
 
-        SubjectExecutionContext.withSubject(_subject, () -> 
message(ChannelMessages.CREATE()));
+        getSubjectExecutionContext().run(() -> 
message(ChannelMessages.CREATE()));
 
         _forceMessageValidation = connection.getContextValue(Boolean.class, 
AMQPConnection_0_8.FORCE_MESSAGE_VALIDATION);
 
@@ -279,7 +278,7 @@ public class AMQChannel extends 
AbstractAMQPSession<AMQChannel, ConsumerTarget_0
 
     public final void receivedComplete()
     {
-        SubjectExecutionContext.withSubject(_subject, this::sync);
+        getSubjectExecutionContext().run(this::sync);
 
     }
 
diff --git 
a/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQPConnection_0_8Impl.java
 
b/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQPConnection_0_8Impl.java
index 5c079b0cf7..f706ccf705 100644
--- 
a/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQPConnection_0_8Impl.java
+++ 
b/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQPConnection_0_8Impl.java
@@ -77,7 +77,6 @@ import 
org.apache.qpid.server.protocol.v0_8.transport.ServerMethodDispatcher;
 import org.apache.qpid.server.protocol.v0_8.transport.ServerMethodProcessor;
 import org.apache.qpid.server.security.AccessDeniedException;
 import org.apache.qpid.server.security.SubjectCreator;
-import org.apache.qpid.server.security.SubjectExecutionContext;
 import org.apache.qpid.server.security.auth.SubjectAuthenticationResult;
 import org.apache.qpid.server.security.auth.sasl.SaslNegotiator;
 import org.apache.qpid.server.session.AMQPSession;
@@ -718,9 +717,11 @@ public class AMQPConnection_0_8Impl
     @Override
     public final void readerIdle()
     {
-        SubjectExecutionContext.withSubject(getSubject(), () -> {
+        runAsSubject(() ->
+        {
             getEventLogger().message(ConnectionMessages.IDLE_CLOSE("Current 
connection state: " + _state, true));
             getNetwork().close();
+            return null;
         });
     }
 
diff --git 
a/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/BrokerDecoder.java
 
b/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/BrokerDecoder.java
index edab027912..5dccce7dbb 100644
--- 
a/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/BrokerDecoder.java
+++ 
b/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/BrokerDecoder.java
@@ -89,7 +89,7 @@ public class BrokerDecoder extends ServerDecoder
             {
                 try
                 {
-                    return 
SubjectExecutionContext.withSubjectUnchecked(channel.getSubject(), () ->
+                    return 
channel.getSubjectExecutionContext().callUnchecked(() ->
                     {
                         int required1;
                         while (true)
diff --git 
a/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/AMQPConnection_1_0Impl.java
 
b/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/AMQPConnection_1_0Impl.java
index 90bdb9d5b8..9602cec9bf 100644
--- 
a/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/AMQPConnection_1_0Impl.java
+++ 
b/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/AMQPConnection_1_0Impl.java
@@ -107,7 +107,6 @@ import 
org.apache.qpid.server.protocol.v1_0.type.transport.Open;
 import org.apache.qpid.server.protocol.v1_0.type.transport.Transfer;
 import org.apache.qpid.server.security.AccessDeniedException;
 import org.apache.qpid.server.security.SubjectCreator;
-import org.apache.qpid.server.security.SubjectExecutionContext;
 import org.apache.qpid.server.security.auth.AuthenticatedPrincipal;
 import org.apache.qpid.server.security.auth.AuthenticationResult;
 import org.apache.qpid.server.security.auth.SubjectAuthenticationResult;
@@ -437,7 +436,7 @@ public class AMQPConnection_1_0Impl extends 
AbstractAMQPConnection<AMQPConnectio
                             : _receivingSessions[frameChannel];
                     if (session != null)
                     {
-                        
SubjectExecutionContext.withSubject(session.getSubject(), () ->
+                        session.getSubjectExecutionContext().run(() ->
                         {
                             ChannelFrameBody channelFrame = channelFrameBody;
                             boolean nextIsSameChannel;
@@ -549,7 +548,7 @@ public class AMQPConnection_1_0Impl extends 
AbstractAMQPConnection<AMQPConnectio
 
         for (final Session_1_0 session : sessions)
         {
-            SubjectExecutionContext.withSubject(session.getSubject(), () -> 
session.remoteEnd(new End()));
+            session.getSubjectExecutionContext().run(() -> 
session.remoteEnd(new End()));
         }
     }
 
@@ -1326,7 +1325,7 @@ public class AMQPConnection_1_0Impl extends 
AbstractAMQPConnection<AMQPConnectio
             {
                 if (session != null)
                 {
-                    SubjectExecutionContext.withSubject(session.getSubject(), 
() -> session.receivedComplete());
+                    
session.getSubjectExecutionContext().run(session::receivedComplete);
                 }
             }
         }
@@ -1578,9 +1577,8 @@ public class AMQPConnection_1_0Impl extends 
AbstractAMQPConnection<AMQPConnectio
                 cause = AmqpError.INTERNAL_ERROR;
         }
         final Session_1_0 actualSession = (Session_1_0) session;
-        addAsyncTask(object ->
-                SubjectExecutionContext.withSubject(actualSession.getSubject(),
-                                                    () -> 
actualSession.close(cause, message)));
+        addAsyncTask(object -> actualSession.getSubjectExecutionContext().run(
+                () -> actualSession.close(cause, message)));
 
     }
 
diff --git 
a/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java
 
b/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java
index fd432b5cca..37e96b53bd 100644
--- 
a/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java
+++ 
b/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java
@@ -195,7 +195,7 @@ public class Session_1_0 extends 
AbstractAMQPSession<Session_1_0, ConsumerTarget
                                       _echoFlowExecutor,
                                       this::sendFlow);
 
-        SubjectExecutionContext.withSubject(_subject, () -> 
_connection.getEventLogger().message(ChannelMessages.CREATE()));
+        getSubjectExecutionContext().run(() -> 
_connection.getEventLogger().message(ChannelMessages.CREATE()));
     }
 
     public void sendDetach(final Detach detach)


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to