Class DurableTopicSubscription

All Implemented Interfaces:
Subscription, SubscriptionRecovery, org.apache.activemq.usage.UsageListener

public class DurableTopicSubscription extends PrefetchSubscription implements org.apache.activemq.usage.UsageListener
  • Constructor Details

    • DurableTopicSubscription

      public DurableTopicSubscription(Broker broker, SystemUsage usageManager, ConnectionContext context, org.apache.activemq.command.ConsumerInfo info, boolean keepDurableSubsActive) throws jakarta.jms.JMSException
      Throws:
      jakarta.jms.JMSException
  • Method Details

    • isActive

      public final boolean isActive()
    • getOfflineTimestamp

      public final long getOfflineTimestamp()
    • setOfflineTimestamp

      public void setOfflineTimestamp(long timestamp)
    • isFull

      public boolean isFull()
      Description copied from class: PrefetchSubscription
      Used to determine if the broker can dispatch to the consumer.
      Specified by:
      isFull in interface Subscription
      Overrides:
      isFull in class PrefetchSubscription
      Returns:
      true if the subscription is full
    • gc

      public void gc()
      Description copied from interface: Subscription
      The subscription should release as may references as it can to help the garbage collector reclaim memory.
      Specified by:
      gc in interface Subscription
      Overrides:
      gc in class AbstractSubscription
    • unmatched

      public void unmatched(org.apache.activemq.broker.region.MessageReference node) throws IOException
      store will have a pending ack for all durables, irrespective of the selector so we need to ack if node is un-matched
      Specified by:
      unmatched in interface Subscription
      Overrides:
      unmatched in class AbstractSubscription
      Throws:
      IOException
    • setPendingBatchSize

      protected void setPendingBatchSize(PendingMessageCursor pending, int numberToDispatch)
      Overrides:
      setPendingBatchSize in class PrefetchSubscription
    • add

      public void add(ConnectionContext context, Destination destination) throws Exception
      Description copied from interface: Subscription
      The subscription will be receiving messages from the destination.
      Specified by:
      add in interface Subscription
      Overrides:
      add in class PrefetchSubscription
      Parameters:
      context -
      destination -
      Throws:
      Exception
    • isEmpty

      public boolean isEmpty(Topic topic)
    • activate

      public void activate(SystemUsage memoryManager, ConnectionContext context, org.apache.activemq.command.ConsumerInfo info, RegionBroker regionBroker) throws Exception
      Throws:
      Exception
    • deactivate

      public void deactivate(boolean keepDurableSubsActive, long lastDeliveredSequenceId) throws Exception
      Throws:
      Exception
    • createMessageDispatch

      protected org.apache.activemq.command.MessageDispatch createMessageDispatch(org.apache.activemq.broker.region.MessageReference node, org.apache.activemq.command.Message message)
      Overrides:
      createMessageDispatch in class PrefetchSubscription
      Parameters:
      node -
      message -
      Returns:
      MessageDispatch
    • add

      public void add(org.apache.activemq.broker.region.MessageReference node) throws Exception
      Description copied from interface: Subscription
      Used to add messages that match the subscription.
      Specified by:
      add in interface Subscription
      Overrides:
      add in class PrefetchSubscription
      Parameters:
      node -
      Throws:
      Exception
    • dispatchPending

      public void dispatchPending() throws IOException
      Overrides:
      dispatchPending in class PrefetchSubscription
      Throws:
      IOException
    • removePending

      public void removePending(org.apache.activemq.broker.region.MessageReference node) throws IOException
      Throws:
      IOException
    • processExpiredAck

      protected void processExpiredAck(ConnectionContext context, Destination dest, org.apache.activemq.broker.region.MessageReference node)
      Overrides:
      processExpiredAck in class PrefetchSubscription
    • doAddRecoveredMessage

      protected void doAddRecoveredMessage(org.apache.activemq.broker.region.MessageReference message) throws Exception
      Overrides:
      doAddRecoveredMessage in class AbstractSubscription
      Throws:
      Exception
    • getPendingQueueSize

      public int getPendingQueueSize()
      Specified by:
      getPendingQueueSize in interface Subscription
      Overrides:
      getPendingQueueSize in class PrefetchSubscription
      Returns:
      number of messages pending delivery
    • setSelector

      public void setSelector(String selector) throws jakarta.jms.InvalidSelectorException
      Description copied from interface: Subscription
      Attempts to change the current active selector on the subscription. This operation is not supported for persistent topics.
      Specified by:
      setSelector in interface Subscription
      Overrides:
      setSelector in class AbstractSubscription
      Throws:
      jakarta.jms.InvalidSelectorException
    • canDispatch

      protected boolean canDispatch(org.apache.activemq.broker.region.MessageReference node)
      Description copied from class: PrefetchSubscription
      Use when a matched message is about to be dispatched to the client.
      Specified by:
      canDispatch in class PrefetchSubscription
      Parameters:
      node -
      Returns:
      false if the message should not be dispatched to the client (another sub may have already dispatched it for example).
    • trackedInPendingTransaction

      protected boolean trackedInPendingTransaction(org.apache.activemq.broker.region.MessageReference node)
      Overrides:
      trackedInPendingTransaction in class PrefetchSubscription
    • acknowledge

      protected void acknowledge(ConnectionContext context, org.apache.activemq.command.MessageAck ack, org.apache.activemq.broker.region.MessageReference node) throws IOException
      Description copied from class: PrefetchSubscription
      Used during acknowledgment to remove the message.
      Specified by:
      acknowledge in class PrefetchSubscription
      Throws:
      IOException
    • toString

      public String toString()
      Overrides:
      toString in class Object
    • getSubscriptionKey

      public SubscriptionKey getSubscriptionKey()
    • destroy

      public void destroy()
      Release any references that we are holding.
      Specified by:
      destroy in interface Subscription
    • onUsageChanged

      public void onUsageChanged(org.apache.activemq.usage.Usage usage, int oldPercentUsage, int newPercentUsage)
      Specified by:
      onUsageChanged in interface org.apache.activemq.usage.UsageListener
    • isDropped

      protected boolean isDropped(org.apache.activemq.broker.region.MessageReference node)
      Specified by:
      isDropped in class PrefetchSubscription
    • isKeepDurableSubsActive

      public boolean isKeepDurableSubsActive()
    • isEnableMessageExpirationOnActiveDurableSubs

      public boolean isEnableMessageExpirationOnActiveDurableSubs()