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]