Package org.ow2.joram.mom.amqp
Class Queue
- java.lang.Object
-
- org.ow2.joram.mom.amqp.Queue
-
- All Implemented Interfaces:
Externalizable,Serializable,QueueMBean
public class Queue extends Object implements QueueMBean, Externalizable
An AMQP queue.- See Also:
- Serialized Form
-
-
Nested Class Summary
Nested Classes Modifier and Type Class Description private static classQueue.Subscriptionprivate static classQueue.SubscriptionKey
-
Field Summary
Fields Modifier and Type Field Description private booleanautodeleteprivate List<String>boundExchangesprivate Map<Queue.SubscriptionKey,Queue.Subscription>consumersprivate booleandurableprivate booleanexclusivestatic longFIRST_DELIVERYstatic Loggerloggerprivate longmsgCounterprivate Stringnameprivate static StringPREFIX_BOUND_EXCHANGEprivate static StringPREFIX_MSGstatic StringPREFIX_QUEUEprivate StringprefixBEprivate StringprefixMsgprivate longproxyIdprivate static longserialVersionUIDdefine serialVersionUID for interoperabilityprivate shortserverIdprivate SortedSet<Message>toAckprivate SortedSet<Message>toDeliver
-
Method Summary
All Methods Static Methods Instance Methods Concrete Methods Modifier and Type Method Description voidackMessages(List<Long> idsToAck)voidaddBoundExchange(String exchange, short serverId, long proxyId)voidcancel(String consumerTag, int channelNumber, short serverId, long proxyId)voidcleanConsumers(short sid)intclear(short serverId, long proxyId)voidconsume(DeliveryListener proxy, int channelId, String consumerTag, boolean exclusiveConsumer, boolean noAck, boolean noLocal, short serverId, long proxyId)private voiddeleteAllMessage(Set<Message> messages)private voiddeleteBoundExchange(String exchangeName)private voiddeleteMessage(long msgId)voiddeleteQueue(String queueName, short serverId, long proxyId)List<String>getBoundExchanges()intgetConsumerCount()List<Deliver>getDeliveries(String consumerTag, int channelId, int maxMessage, short serverId, long proxyId)longgetHandledMessageCount()AMQP.Queue.DeclareOkgetInfo(short serverId, long proxyId)StringgetName()intgetToAckMessageCount()intgetToDeliverMessageCount()booleanisAutodelete()booleanisDurable()booleanisExclusive()static QueueloadQueue(String name)voidpublish(Message msg, boolean immediate, short serverId, long proxyId)voidreadExternal(ObjectInput in)Messagereceive(boolean noAck, short serverId, long proxyId)voidrecoverMessages(List<Long> idsToRecover)voidremoveBoundExchange(String exchangeName)voidremoveBoundExchange(String exchangeName, short serverId, long proxyId)private voidsaveBoundExchange(String exchange)private voidsaveMessage(Message msg)private voidsaveQueue(Queue queue)voidwriteExternal(ObjectOutput out)
-
-
-
Field Detail
-
logger
public static final Logger logger
-
serialVersionUID
private static final long serialVersionUID
define serialVersionUID for interoperability- See Also:
- Constant Field Values
-
FIRST_DELIVERY
public static final long FIRST_DELIVERY
- See Also:
- Constant Field Values
-
name
private String name
-
durable
private boolean durable
-
autodelete
private boolean autodelete
-
exclusive
private boolean exclusive
-
serverId
private short serverId
-
proxyId
private long proxyId
-
msgCounter
private long msgCounter
-
consumers
private Map<Queue.SubscriptionKey,Queue.Subscription> consumers
-
PREFIX_QUEUE
public static final String PREFIX_QUEUE
- See Also:
- Constant Field Values
-
PREFIX_MSG
private static final String PREFIX_MSG
- See Also:
- Constant Field Values
-
PREFIX_BOUND_EXCHANGE
private static final String PREFIX_BOUND_EXCHANGE
- See Also:
- Constant Field Values
-
prefixMsg
private String prefixMsg
-
prefixBE
private String prefixBE
-
-
Constructor Detail
-
Queue
public Queue()
-
Queue
public Queue(String name, boolean durable, boolean autodelete, boolean exclusive, short serverId, long proxyId) throws TransactionException
- Throws:
TransactionException
-
-
Method Detail
-
receive
public Message receive(boolean noAck, short serverId, long proxyId) throws ResourceLockedException, TransactionException
-
consume
public void consume(DeliveryListener proxy, int channelId, String consumerTag, boolean exclusiveConsumer, boolean noAck, boolean noLocal, short serverId, long proxyId) throws AccessRefusedException, ResourceLockedException
-
getDeliveries
public List<Deliver> getDeliveries(String consumerTag, int channelId, int maxMessage, short serverId, long proxyId)
-
publish
public void publish(Message msg, boolean immediate, short serverId, long proxyId) throws NoConsumersException, TransactionException
-
cancel
public void cancel(String consumerTag, int channelNumber, short serverId, long proxyId) throws ResourceLockedException
- Throws:
ResourceLockedException
-
cleanConsumers
public void cleanConsumers(short sid)
-
clear
public int clear(short serverId, long proxyId) throws ResourceLockedException, TransactionException
-
recoverMessages
public void recoverMessages(List<Long> idsToRecover) throws TransactionException
- Throws:
TransactionException
-
getInfo
public AMQP.Queue.DeclareOk getInfo(short serverId, long proxyId) throws ResourceLockedException
- Throws:
ResourceLockedException
-
getName
public String getName()
- Specified by:
getNamein interfaceQueueMBean
-
getConsumerCount
public int getConsumerCount()
- Specified by:
getConsumerCountin interfaceQueueMBean
-
isAutodelete
public boolean isAutodelete()
- Specified by:
isAutodeletein interfaceQueueMBean
-
getToDeliverMessageCount
public int getToDeliverMessageCount()
- Specified by:
getToDeliverMessageCountin interfaceQueueMBean
-
getToAckMessageCount
public int getToAckMessageCount()
- Specified by:
getToAckMessageCountin interfaceQueueMBean
-
getHandledMessageCount
public long getHandledMessageCount()
- Specified by:
getHandledMessageCountin interfaceQueueMBean
-
getBoundExchanges
public List<String> getBoundExchanges()
- Specified by:
getBoundExchangesin interfaceQueueMBean
-
isDurable
public boolean isDurable()
- Specified by:
isDurablein interfaceQueueMBean
-
isExclusive
public boolean isExclusive()
- Specified by:
isExclusivein interfaceQueueMBean
-
addBoundExchange
public void addBoundExchange(String exchange, short serverId, long proxyId) throws TransactionException, ResourceLockedException
-
removeBoundExchange
public void removeBoundExchange(String exchangeName)
-
removeBoundExchange
public void removeBoundExchange(String exchangeName, short serverId, long proxyId) throws ResourceLockedException
- Throws:
ResourceLockedException
-
deleteQueue
public void deleteQueue(String queueName, short serverId, long proxyId) throws ResourceLockedException, TransactionException
-
loadQueue
public static Queue loadQueue(String name) throws IOException, ClassNotFoundException, TransactionException
-
saveQueue
private void saveQueue(Queue queue) throws TransactionException
- Throws:
TransactionException
-
saveBoundExchange
private void saveBoundExchange(String exchange) throws TransactionException
- Throws:
TransactionException
-
deleteBoundExchange
private void deleteBoundExchange(String exchangeName)
-
saveMessage
private void saveMessage(Message msg) throws TransactionException
- Throws:
TransactionException
-
deleteMessage
private void deleteMessage(long msgId)
-
writeExternal
public void writeExternal(ObjectOutput out) throws IOException
- Specified by:
writeExternalin interfaceExternalizable- Parameters:
out-- Throws:
IOException
-
readExternal
public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException
- Specified by:
readExternalin interfaceExternalizable- Parameters:
in-- Throws:
IOExceptionClassNotFoundException
-
-