summaryrefslogtreecommitdiff
path: root/qpid/cpp/src/qpid/client/SubscriptionManager.h
blob: 07faa48fee1ff37da10154bfa4c9904e91e8a350 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
#ifndef QPID_CLIENT_SUBSCRIPTIONMANAGER_H
#define QPID_CLIENT_SUBSCRIPTIONMANAGER_H

/*
 *
 * Licensed to the Apache Software Foundation (ASF) under one
 * or more contributor license agreements.  See the NOTICE file
 * distributed with this work for additional information
 * regarding copyright ownership.  The ASF licenses this file
 * to you under the Apache License, Version 2.0 (the
 * "License"); you may not use this file except in compliance
 * with the License.  You may obtain a copy of the License at
 * 
 *   http://www.apache.org/licenses/LICENSE-2.0
 * 
 * Unless required by applicable law or agreed to in writing,
 * software distributed under the License is distributed on an
 * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
 * KIND, either express or implied.  See the License for the
 * specific language governing permissions and limitations
 * under the License.
 *
 */
#include "qpid/sys/Mutex.h"
#include <qpid/client/Dispatcher.h>
#include <qpid/client/Completion.h>
#include <qpid/client/Session.h>
#include <qpid/client/MessageListener.h>
#include <qpid/client/LocalQueue.h>
#include <qpid/client/FlowControl.h>
#include <qpid/sys/Runnable.h>
#include <set>
#include <sstream>

namespace qpid {
namespace client {

/**
 * A class to help create and manage subscriptions.
 * 
 * Set up your subscriptions, then call run() to have messages
 * delivered.
 *  
 * \ingroup clientapi
 */
class SubscriptionManager : public sys::Runnable
{
    typedef sys::Mutex::ScopedLock Lock;
    typedef sys::Mutex::ScopedUnlock Unlock;

    void subscribeInternal(const std::string& q, const std::string& dest, const FlowControl&);
    
    qpid::client::Dispatcher dispatcher;
    qpid::client::AsyncSession session;
    FlowControl flowControl;
    AckPolicy autoAck;
    bool acceptMode;
    bool acquireMode;
    bool autoStop;
    
  public:
    /** Create a new SubscriptionManager associated with a session */
    SubscriptionManager(const Session& session);
    
    /**
     * Subscribe a MessagesListener to receive messages from queue.
     *
     * Provide your own subclass of MessagesListener to process
     * incoming messages. It will be called for each message received.
     * 
     *@param listener Listener object to receive messages.
     *@param queue Name of the queue to subscribe to.
     *@param flow initial FlowControl for the subscription.
     *@param tag Unique destination tag for the listener.
     * If not specified, the queue name is used.
     */
    void subscribe(MessageListener& listener,
                   const std::string& queue,
                   const FlowControl& flow,
                   const std::string& tag=std::string());

    /**
     * Subscribe a LocalQueue to receive messages from queue.
     * 
     * Incoming messages are stored in the queue for you to retrieve.
     * 
     *@param queue Name of the queue to subscribe to.
     *@param flow initial FlowControl for the subscription.
     *@param tag Unique destination tag for the listener.
     * If not specified, the queue name is used.
     */
    void subscribe(LocalQueue& localQueue,
                   const std::string& queue,
                   const FlowControl& flow,
                   const std::string& tag=std::string());

    /**
     * Subscribe a MessagesListener to receive messages from queue.
     *
     * Provide your own subclass of MessagesListener to process
     * incoming messages. It will be called for each message received.
     * 
     *@param listener Listener object to receive messages.
     *@param queue Name of the queue to subscribe to.
     *@param tag Unique destination tag for the listener.
     * If not specified, the queue name is used.
     */
    void subscribe(MessageListener& listener,
                   const std::string& queue,
                   const std::string& tag=std::string());

    /**
     * Subscribe a LocalQueue to receive messages from queue.
     * 
     * Incoming messages are stored in the queue for you to retrieve.
     * 
     *@param queue Name of the queue to subscribe to.
     *@param tag Unique destination tag for the listener.
     * If not specified, the queue name is used.
     */
    void subscribe(LocalQueue& localQueue,
                   const std::string& queue,
                   const std::string& tag=std::string());


    /** Get a single message from a queue.
     *@param result is set to the message from the queue.
     *@
     *@param timeout wait up this timeout for a message to appear. 
     *@return true if result was set, false if no message available after timeout.
     */
    bool get(Message& result, const std::string& queue, sys::Duration timeout=0);

    /** Cancel a subscription. */
    void cancel(const std::string tag);

    /** Deliver messages in the current thread until stop() is called.
     * Only one thread may be running in a SubscriptionManager at a time.
     * @see run
     */
    void run();

    /** Start a new thread to deliver messages.
     * Only one thread may be running in a SubscriptionManager at a time.
     * @see start
     */
    void start();
    
    /** If set true, run() will stop when all subscriptions
     * are cancelled. If false, run will only stop when stop()
     * is called. True by default.
     */
    void setAutoStop(bool set=true);

    /** Stop delivery. Causes run() to return, or the thread started with start() to exit. */
    void stop();

    static const uint32_t UNLIMITED=0xFFFFFFFF;

    /** Set the flow control for destination. */
    void setFlowControl(const std::string& destintion, const FlowControl& flow);

    /** Set the default initial flow control for subscriptions that do not specify it. */
    void setFlowControl(const FlowControl& flow);

    /** Get the default flow control for new subscriptions that do not specify it. */
    const FlowControl& getFlowControl() const;

    /** Set the flow control for destination tag.
     *@param tag: name of the destination.
     *@param messages: message credit.
     *@param bytes: byte credit.
     *@param window: if true use window-based flow control.
     */
    void setFlowControl(const std::string& tag, uint32_t messages,  uint32_t bytes, bool window=true);

    /** Set the initial flow control settings to be applied to each new subscribtion.
     *@param messages: message credit.
     *@param bytes: byte credit.
     *@param window: if true use window-based flow control.
     */
    void setFlowControl(uint32_t messages,  uint32_t bytes, bool window=true);

    /** Set the accept-mode for new subscriptions. Defaults to true.
     *@param required: if true messages must be confirmed by calling
     *Message::acknowledge() or automatically via an AckPolicy, see setAckPolicy()
     */
    void setAcceptMode(bool required);

    /** Set the acquire-mode for new subscriptions. Defaults to false.
     *@param acquire: if false messages pre-acquired, if true
     * messages are dequed on acknowledgement or on transfer 
     * depending on acceptMode.
     */
    void setAcquireMode(bool acquire);

    /** Set the acknowledgement policy for new subscriptions.
     * Default is to acknowledge every message automatically.
     */
    void setAckPolicy(const AckPolicy& autoAck);

    AckPolicy& getAckPolicy();

    void registerFailoverHandler ( boost::function<void ()> fh );

    Session getSession() const;
};

/** AutoCancel cancels a subscription in its destructor */
class AutoCancel {
  public:
    AutoCancel(SubscriptionManager& sm_, const std::string& tag_) : sm(sm_), tag(tag_) {}
    ~AutoCancel() { sm.cancel(tag); }
  private:
    SubscriptionManager& sm;
    std::string tag;
};

}} // namespace qpid::client

#endif  /*!QPID_CLIENT_SUBSCRIPTIONMANAGER_H*/