summaryrefslogtreecommitdiff
path: root/cpp/client/inc/Channel.h
blob: debecf922e1148895bdec054d82adeeade98aaa0 (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
/*
 *
 * Copyright (c) 2006 The Apache Software Foundation
 *
 * Licensed 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 <map>
#include <string>
#include <queue>
#include "sys/types.h"

#ifndef _Channel_
#define _Channel_

#include "amqp_framing.h"

#include "ThreadFactory.h"

#include "Connection.h"
#include "Exchange.h"
#include "IncomingMessage.h"
#include "Message.h"
#include "MessageListener.h"
#include "Queue.h"
#include "ResponseHandler.h"
#include "ReturnedMessageHandler.h"

namespace qpid {
namespace client {
    enum ack_modes {NO_ACK=0, AUTO_ACK=1, LAZY_ACK=2, CLIENT_ACK=3};

    class Channel : private virtual qpid::framing::BodyHandler, public virtual qpid::concurrent::Runnable{
        struct Consumer{
            MessageListener* listener;
            int ackMode;
            int count;
            u_int64_t lastDeliveryTag;
        };
        typedef std::map<string,Consumer*>::iterator consumer_iterator; 

	u_int16_t id;
	Connection* con;
	qpid::concurrent::ThreadFactory* threadFactory;
	qpid::concurrent::Thread* dispatcher;
	qpid::framing::OutputHandler* out;
	IncomingMessage* incoming;
	ResponseHandler responses;
	std::queue<IncomingMessage*> messages;//holds returned messages or those delivered for a consume
	IncomingMessage* retrieved;//holds response to basic.get
	qpid::concurrent::Monitor* dispatchMonitor;
	qpid::concurrent::Monitor* retrievalMonitor;
	std::map<std::string, Consumer*> consumers;
	ReturnedMessageHandler* returnsHandler;
	bool closed;

        u_int16_t prefetch;
        const bool transactional;

	void enqueue();
	void retrieve(Message& msg);
	IncomingMessage* dequeue();
	void dispatch();
	void stop();
	void sendAndReceive(qpid::framing::AMQFrame* frame, const qpid::framing::AMQMethodBody& body);            
        void deliver(Consumer* consumer, Message& msg);
        void setQos();
	void cancelAll();

	virtual void handleMethod(qpid::framing::AMQMethodBody::shared_ptr body);
	virtual void handleHeader(qpid::framing::AMQHeaderBody::shared_ptr body);
	virtual void handleContent(qpid::framing::AMQContentBody::shared_ptr body);
	virtual void handleHeartbeat(qpid::framing::AMQHeartbeatBody::shared_ptr body);

    public:
	Channel(bool transactional = false, u_int16_t prefetch = 500);
	~Channel();

	void declareExchange(Exchange& exchange, bool synch = true);
	void deleteExchange(Exchange& exchange, bool synch = true);
	void declareQueue(Queue& queue, bool synch = true);
	void deleteQueue(Queue& queue, bool ifunused = false, bool ifempty = false, bool synch = true);
	void bind(const Exchange& exchange, const Queue& queue, const std::string& key, 
                  const qpid::framing::FieldTable& args, bool synch = true);
        void consume(Queue& queue, std::string& tag, MessageListener* listener, 
                     int ackMode = NO_ACK, bool noLocal = false, bool synch = true);
	void cancel(std::string& tag, bool synch = true);
        bool get(Message& msg, const Queue& queue, int ackMode = NO_ACK);
        void publish(Message& msg, const Exchange& exchange, const std::string& routingKey, 
                     bool mandatory = false, bool immediate = false);

        void commit();
        void rollback();

        void setPrefetch(u_int16_t prefetch);

	/**
	 * Start message dispatching on a new thread
	 */
	void start();
	/**
	 * Do message dispatching on this thread
	 */
	void run();

        void close();

	void setReturnedMessageHandler(ReturnedMessageHandler* handler);

        friend class Connection;
    };

}
}


#endif