forked from confluentinc/librdkafka
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathQueueImpl.cpp
66 lines (56 loc) · 2.36 KB
/
QueueImpl.cpp
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
/*
* librdkafka - Apache Kafka C/C++ library
*
* Copyright (c) 2014 Magnus Edenhill
* All rights reserved.
*
* Redistribution and use in source and binary forms, with or without
* modification, are permitted provided that the following conditions are met:
*
* 1. Redistributions of source code must retain the above copyright notice,
* this list of conditions and the following disclaimer.
* 2. Redistributions in binary form must reproduce the above copyright notice,
* this list of conditions and the following disclaimer in the documentation
* and/or other materials provided with the distribution.
*
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS"
* AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE
* IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE
* ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT OWNER OR CONTRIBUTORS BE
* LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR
* CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF
* SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS
* INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN
* CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE)
* ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE
* POSSIBILITY OF SUCH DAMAGE.
*/
#include <cerrno>
#include "rdkafkacpp_int.h"
RdKafka::Queue::~Queue () {
}
RdKafka::Queue *RdKafka::Queue::create (Handle *base) {
RdKafka::QueueImpl *queueimpl = new RdKafka::QueueImpl;
queueimpl->queue_ = rd_kafka_queue_new(dynamic_cast<HandleImpl*>(base)->rk_);
return queueimpl;
}
RdKafka::ErrorCode
RdKafka::QueueImpl::forward (Queue *queue) {
if (!queue) {
rd_kafka_queue_forward(queue_, NULL);
} else {
QueueImpl *queueimpl = dynamic_cast<QueueImpl *>(queue);
rd_kafka_queue_forward(queue_, queueimpl->queue_);
}
return RdKafka::ERR_NO_ERROR;
}
RdKafka::Message *RdKafka::QueueImpl::consume (int timeout_ms) {
rd_kafka_message_t *rkmessage;
rkmessage = rd_kafka_consume_queue(queue_, timeout_ms);
if (!rkmessage)
return new RdKafka::MessageImpl(NULL, RdKafka::ERR__TIMED_OUT);
return new RdKafka::MessageImpl(rkmessage);
}
int RdKafka::QueueImpl::poll (int timeout_ms) {
return rd_kafka_queue_poll_callback(queue_, timeout_ms);
}