Class TransportConnection
java.lang.Object
org.apache.activemq.broker.TransportConnection
- All Implemented Interfaces:
Connection, org.apache.activemq.Service, org.apache.activemq.state.CommandVisitor, org.apache.activemq.thread.Task
- Direct Known Subclasses:
ManagedTransportConnection
public class TransportConnection
extends Object
implements Connection, org.apache.activemq.thread.Task, org.apache.activemq.state.CommandVisitor
-
Field Summary
FieldsModifier and TypeFieldDescriptionprotected final Brokerprotected final Map<org.apache.activemq.command.ConnectionId, org.apache.activemq.state.ConnectionState> protected org.apache.activemq.command.BrokerInfoprotected final BrokerServiceprotected final TransportConnectorprotected final List<org.apache.activemq.command.Command> protected AtomicBooleanprotected org.apache.activemq.thread.TaskRunnerprotected final AtomicReference<Throwable> -
Constructor Summary
ConstructorsConstructorDescriptionTransportConnection(TransportConnector connector, org.apache.activemq.transport.Transport transport, Broker broker, org.apache.activemq.thread.TaskRunnerFactory taskRunnerFactory, org.apache.activemq.thread.TaskRunnerFactory stopTaskRunnerFactory) -
Method Summary
Modifier and TypeMethodDescriptionvoiddelayedStop(int waitTime, String reason, Throwable cause) protected voiddispatch(org.apache.activemq.command.Command command) voiddispatchAsync(org.apache.activemq.command.Command message) Sends a message to the client.voiddispatchSync(org.apache.activemq.command.Command message) Sends a message to the client.voiddoMark()Mark the Connection, so we can deem if it's collectable on the next sweepprotected voiddoStop()intReturns the number of active transactions established on this Connection.Returns the time in ms since epoch when connection was established.intReturns the number of messages to be dispatched to this connectionReturns the number of active transactions established on this Connection.getProducerBrokerExchangeIfExists(org.apache.activemq.command.ProducerInfo producerInfo) intorg.apache.activemq.command.WireFormatInfoReturns the statistics for this connectionorg.apache.activemq.transport.TransportbooleanisActive()booleanbooleanbooleanbooleanbooleanbooleanbooleanbooleanreturn true if a network connectionbooleanbooleanisSlow()booleanbooleanbooleaniterate()protected List<TransportConnectionState> protected TransportConnectionStatelookupConnectionState(String connectionId) lookupConnectionState(org.apache.activemq.command.ConnectionId connectionId) protected TransportConnectionStatelookupConnectionState(org.apache.activemq.command.ConsumerId id) protected TransportConnectionStatelookupConnectionState(org.apache.activemq.command.ProducerId id) protected TransportConnectionStatelookupConnectionState(org.apache.activemq.command.SessionId id) org.apache.activemq.command.ResponseprocessAddConnection(org.apache.activemq.command.ConnectionInfo info) org.apache.activemq.command.ResponseprocessAddConsumer(org.apache.activemq.command.ConsumerInfo info) org.apache.activemq.command.ResponseprocessAddDestination(org.apache.activemq.command.DestinationInfo info) org.apache.activemq.command.ResponseprocessAddProducer(org.apache.activemq.command.ProducerInfo info) org.apache.activemq.command.ResponseprocessAddSession(org.apache.activemq.command.SessionInfo info) org.apache.activemq.command.ResponseprocessBeginTransaction(org.apache.activemq.command.TransactionInfo info) org.apache.activemq.command.ResponseprocessBrokerInfo(org.apache.activemq.command.BrokerInfo info) org.apache.activemq.command.ResponseprocessBrokerSubscriptionInfo(org.apache.activemq.command.BrokerSubscriptionInfo info) org.apache.activemq.command.ResponseprocessCommitTransactionOnePhase(org.apache.activemq.command.TransactionInfo info) org.apache.activemq.command.ResponseprocessCommitTransactionTwoPhase(org.apache.activemq.command.TransactionInfo info) org.apache.activemq.command.ResponseprocessConnectionControl(org.apache.activemq.command.ConnectionControl control) org.apache.activemq.command.ResponseprocessConnectionError(org.apache.activemq.command.ConnectionError error) org.apache.activemq.command.ResponseprocessConsumerControl(org.apache.activemq.command.ConsumerControl control) org.apache.activemq.command.ResponseprocessControlCommand(org.apache.activemq.command.ControlCommand command) protected voidprocessDispatch(org.apache.activemq.command.Command command) org.apache.activemq.command.ResponseprocessEndTransaction(org.apache.activemq.command.TransactionInfo info) org.apache.activemq.command.ResponseprocessFlush(org.apache.activemq.command.FlushCommand command) org.apache.activemq.command.ResponseprocessForgetTransaction(org.apache.activemq.command.TransactionInfo info) org.apache.activemq.command.ResponseprocessKeepAlive(org.apache.activemq.command.KeepAliveInfo info) org.apache.activemq.command.ResponseprocessMessage(org.apache.activemq.command.Message messageSend) org.apache.activemq.command.ResponseprocessMessageAck(org.apache.activemq.command.MessageAck ack) org.apache.activemq.command.ResponseprocessMessageDispatch(org.apache.activemq.command.MessageDispatch dispatch) org.apache.activemq.command.ResponseprocessMessageDispatchNotification(org.apache.activemq.command.MessageDispatchNotification notification) org.apache.activemq.command.ResponseprocessMessagePull(org.apache.activemq.command.MessagePull pull) org.apache.activemq.command.ResponseprocessPrepareTransaction(org.apache.activemq.command.TransactionInfo info) org.apache.activemq.command.ResponseprocessProducerAck(org.apache.activemq.command.ProducerAck ack) org.apache.activemq.command.ResponseprocessRecoverTransactions(org.apache.activemq.command.TransactionInfo info) org.apache.activemq.command.ResponseprocessRemoveConnection(org.apache.activemq.command.ConnectionId id, long lastDeliveredSequenceId) org.apache.activemq.command.ResponseprocessRemoveConsumer(org.apache.activemq.command.ConsumerId id, long lastDeliveredSequenceId) org.apache.activemq.command.ResponseprocessRemoveDestination(org.apache.activemq.command.DestinationInfo info) org.apache.activemq.command.ResponseprocessRemoveProducer(org.apache.activemq.command.ProducerId id) org.apache.activemq.command.ResponseprocessRemoveSession(org.apache.activemq.command.SessionId id, long lastDeliveredSequenceId) org.apache.activemq.command.ResponseprocessRemoveSubscription(org.apache.activemq.command.RemoveSubscriptionInfo info) org.apache.activemq.command.ResponseprocessRollbackTransaction(org.apache.activemq.command.TransactionInfo info) org.apache.activemq.command.ResponseprocessShutdown(org.apache.activemq.command.ShutdownInfo info) org.apache.activemq.command.ResponseprocessWireFormat(org.apache.activemq.command.WireFormatInfo info) protected TransportConnectionStateregisterConnectionState(org.apache.activemq.command.ConnectionId connectionId, TransportConnectionState state) org.apache.activemq.command.Responseservice(org.apache.activemq.command.Command command) Services a client command and submits it to the broker.voidCloses a clients connection due to a detected error.voidCalls the serviceException method in an async thread.voidvoidsetActive(boolean active) voidsetBlocked(boolean blocked) voidsetBlockedCandidate(boolean blockedCandidate) voidsetConnected(boolean connected) protected voidsetDuplexNetworkConnectorId(String duplexNetworkConnectorId) voidsetMarkedCandidate(boolean markedCandidate) voidsetMessageAuthorizationPolicy(MessageAuthorizationPolicy messageAuthorizationPolicy) voidsetSlow(boolean slow) voidstart()voidstop()voidvoidtoString()protected TransportConnectionStateunregisterConnectionState(org.apache.activemq.command.ConnectionId connectionId) voidupdateClient(org.apache.activemq.command.ConnectionControl control)
-
Field Details
-
broker
-
brokerService
-
connector
-
brokerConnectionStates
protected final Map<org.apache.activemq.command.ConnectionId, org.apache.activemq.state.ConnectionState> brokerConnectionStates -
brokerInfo
protected org.apache.activemq.command.BrokerInfo brokerInfo -
dispatchQueue
-
taskRunner
protected org.apache.activemq.thread.TaskRunner taskRunner -
transportException
-
dispatchStopped
-
-
Constructor Details
-
TransportConnection
public TransportConnection(TransportConnector connector, org.apache.activemq.transport.Transport transport, Broker broker, org.apache.activemq.thread.TaskRunnerFactory taskRunnerFactory, org.apache.activemq.thread.TaskRunnerFactory stopTaskRunnerFactory) - Parameters:
taskRunnerFactory- - can be null if you want direct dispatch to the transport else commands are sent async.stopTaskRunnerFactory- - can not be null, used for stopping this connection.
-
-
Method Details
-
getDispatchQueueSize
public int getDispatchQueueSize()Returns the number of messages to be dispatched to this connection- Specified by:
getDispatchQueueSizein interfaceConnection- Returns:
- size of dispatch queue
-
serviceTransportException
-
serviceExceptionAsync
Calls the serviceException method in an async thread. Since handling a service exception closes a socket, we should not tie up broker threads since client sockets may hang or cause deadlocks.- Specified by:
serviceExceptionAsyncin interfaceConnection
-
serviceException
Closes a clients connection due to a detected error. Errors are ignored if: the client is closing or broker is closing. Otherwise, the connection error transmitted to the client before stopping it's transport.- Specified by:
serviceExceptionin interfaceConnection- Parameters:
e-
-
service
public org.apache.activemq.command.Response service(org.apache.activemq.command.Command command) Description copied from interface:ConnectionServices a client command and submits it to the broker.- Specified by:
servicein interfaceConnection- Parameters:
command-- Returns:
- Response
-
processKeepAlive
-
processRemoveSubscription
-
processWireFormat
-
processShutdown
-
processFlush
-
processBeginTransaction
-
getActiveTransactionCount
public int getActiveTransactionCount()Description copied from interface:ConnectionReturns the number of active transactions established on this Connection.- Specified by:
getActiveTransactionCountin interfaceConnection- Returns:
- the number of active transactions established on this Connection..
-
getOldestActiveTransactionDuration
Description copied from interface:ConnectionReturns the number of active transactions established on this Connection.- Specified by:
getOldestActiveTransactionDurationin interfaceConnection- Returns:
- the number of active transactions established on this Connection..
-
processEndTransaction
-
processPrepareTransaction
-
processCommitTransactionOnePhase
-
processCommitTransactionTwoPhase
-
processRollbackTransaction
-
processForgetTransaction
-
processRecoverTransactions
-
processMessage
-
processMessageAck
-
processMessagePull
-
processMessageDispatchNotification
-
processAddDestination
-
processRemoveDestination
-
processAddProducer
-
processRemoveProducer
-
processAddConsumer
-
processRemoveConsumer
-
processAddSession
-
processRemoveSession
-
processAddConnection
-
processRemoveConnection
public org.apache.activemq.command.Response processRemoveConnection(org.apache.activemq.command.ConnectionId id, long lastDeliveredSequenceId) throws InterruptedException - Specified by:
processRemoveConnectionin interfaceorg.apache.activemq.state.CommandVisitor- Throws:
InterruptedException
-
processProducerAck
-
getConnector
- Specified by:
getConnectorin interfaceConnection- Returns:
- the connector that created this connection.
-
dispatchSync
public void dispatchSync(org.apache.activemq.command.Command message) Description copied from interface:ConnectionSends a message to the client.- Specified by:
dispatchSyncin interfaceConnection- Parameters:
message- the message to send to the client.
-
dispatchAsync
public void dispatchAsync(org.apache.activemq.command.Command message) Description copied from interface:ConnectionSends a message to the client.- Specified by:
dispatchAsyncin interfaceConnection- Parameters:
message-
-
processDispatch
- Throws:
IOException
-
iterate
public boolean iterate()- Specified by:
iteratein interfaceorg.apache.activemq.thread.Task
-
getStatistics
Returns the statistics for this connection- Specified by:
getStatisticsin interfaceConnection
-
getMessageAuthorizationPolicy
-
setMessageAuthorizationPolicy
-
isManageable
public boolean isManageable()- Specified by:
isManageablein interfaceConnection- Returns:
- true if the Connection will process control commands
-
start
-
stop
-
delayedStop
-
stopAsync
-
stopAsync
public void stopAsync() -
toString
-
doStop
-
isBlockedCandidate
public boolean isBlockedCandidate()- Returns:
- Returns the blockedCandidate.
-
setBlockedCandidate
public void setBlockedCandidate(boolean blockedCandidate) - Parameters:
blockedCandidate- The blockedCandidate to set.
-
isMarkedCandidate
public boolean isMarkedCandidate()- Returns:
- Returns the markedCandidate.
-
setMarkedCandidate
public void setMarkedCandidate(boolean markedCandidate) - Parameters:
markedCandidate- The markedCandidate to set.
-
setSlow
public void setSlow(boolean slow) - Parameters:
slow- The slow to set.
-
isSlow
public boolean isSlow()- Specified by:
isSlowin interfaceConnection- Returns:
- true if the Connection is slow
-
isMarkedBlockedCandidate
public boolean isMarkedBlockedCandidate()- Returns:
- true if the Connection is potentially blocked
-
doMark
public void doMark()Mark the Connection, so we can deem if it's collectable on the next sweep -
isBlocked
public boolean isBlocked()- Specified by:
isBlockedin interfaceConnection- Returns:
- if after being marked, the Connection is still writing
-
isConnected
public boolean isConnected()- Specified by:
isConnectedin interfaceConnection- Returns:
- true if the Connection is connected
-
setBlocked
public void setBlocked(boolean blocked) - Parameters:
blocked- The blocked to set.
-
setConnected
public void setConnected(boolean connected) - Parameters:
connected- The connected to set.
-
isActive
public boolean isActive()- Specified by:
isActivein interfaceConnection- Returns:
- true if the Connection is active
-
setActive
public void setActive(boolean active) - Parameters:
active- The active to set.
-
isStarting
public boolean isStarting()- Returns:
- true if the Connection is starting
-
isNetworkConnection
public boolean isNetworkConnection()Description copied from interface:Connectionreturn true if a network connection- Specified by:
isNetworkConnectionin interfaceConnection- Returns:
- if this is a network connection
-
isFaultTolerantConnection
public boolean isFaultTolerantConnection()- Specified by:
isFaultTolerantConnectionin interfaceConnection- Returns:
- true if a fault tolerant connection
-
isPendingStop
public boolean isPendingStop()- Returns:
- true if the Connection needs to stop
-
processBrokerInfo
public org.apache.activemq.command.Response processBrokerInfo(org.apache.activemq.command.BrokerInfo info) throws IOException - Specified by:
processBrokerInfoin interfaceorg.apache.activemq.state.CommandVisitor- Throws:
IOException
-
dispatch
- Throws:
IOException
-
getRemoteAddress
- Specified by:
getRemoteAddressin interfaceConnection- Returns:
- the source address for this connection
-
getTransport
public org.apache.activemq.transport.Transport getTransport() -
getConnectionId
- Specified by:
getConnectionIdin interfaceConnection
-
updateClient
public void updateClient(org.apache.activemq.command.ConnectionControl control) - Specified by:
updateClientin interfaceConnection
-
getProducerBrokerExchangeIfExists
public ProducerBrokerExchange getProducerBrokerExchangeIfExists(org.apache.activemq.command.ProducerInfo producerInfo) -
getProtocolVersion
public int getProtocolVersion() -
processControlCommand
-
processMessageDispatch
-
processConnectionControl
-
processConnectionError
-
processConsumerControl
-
registerConnectionState
protected TransportConnectionState registerConnectionState(org.apache.activemq.command.ConnectionId connectionId, TransportConnectionState state) -
unregisterConnectionState
protected TransportConnectionState unregisterConnectionState(org.apache.activemq.command.ConnectionId connectionId) -
listConnectionStates
-
lookupConnectionState
-
lookupConnectionState
-
lookupConnectionState
-
lookupConnectionState
-
lookupConnectionState
public TransportConnectionState lookupConnectionState(org.apache.activemq.command.ConnectionId connectionId) -
setDuplexNetworkConnectorId
-
getDuplexNetworkConnectorId
-
isStopping
public boolean isStopping() -
getStopped
-
getRemoteWireFormatInfo
public org.apache.activemq.command.WireFormatInfo getRemoteWireFormatInfo() -
processBrokerSubscriptionInfo
-
getConnectedTimestamp
Description copied from interface:ConnectionReturns the time in ms since epoch when connection was established.- Specified by:
getConnectedTimestampin interfaceConnection- Returns:
- time in ms since epoch when connection was established.
-