FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

xpending rework · cpp-redis/cpp_redis@7e50173 · GitHub

Repository navigation

Commit 7e50173

Browse files
committed
xpending rework
added overload methods for xpending added range class for simplifying range arguments (reduces number of overloads)
1 parent 61b5659 commit 7e50173

4 files changed

Lines changed: 249 additions & 24 deletions

File tree

‎includes/cpp_redis/core/client.hpp‎

Lines changed: 37 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232
#include <string>
3333
#include <vector>
3434

35+
#include <cpp_redis/core/types.hpp>
3536
#include <cpp_redis/core/sentinel.hpp>
3637
#include <cpp_redis/helpers/variadic_template.hpp>
3738
#include <cpp_redis/misc/logger.hpp>
@@ -1563,23 +1564,53 @@ namespace cpp_redis {
15631564

15641565
client &xpending(const std::string &key,
15651566
const std::string &group_name,
1566-
const int &start = 0,
1567-
const int &end = -1,
1568-
const int &count = 10,
1569-
const std::string &consumer_name = "",
1567+
const reply_callback_t &reply_callback);
1568+
1569+
client &xpending(const std::string &key,
1570+
const std::string &group_name,
1571+
const range_t &range,
1572+
const reply_callback_t &reply_callback);
1573+
1574+
client &xpending(const std::string &key,
1575+
const std::string &group_name,
1576+
const std::string &consumer_name,
1577+
const reply_callback_t &reply_callback);
1578+
1579+
client &xpending(const std::string &key,
1580+
const std::string &group_name,
1581+
const range_t &range,
1582+
const std::string &consumer_name,
15701583
const reply_callback_t &reply_callback = nullptr);
15711584

15721585
//XREADGROUP GROUP group consumer [COUNT count] [BLOCK milliseconds] [NOACK] STREAMS key [key ...] ID [ID ...]
15731586

15741587
client &xreadgroup(const std::string &group_name,
15751588
const std::string &consumer_name,
15761589
int count,
1577-
int block_milli_sec,
15781590
bool no_ack,
15791591
std::multimap<std::string, std::string> stream_members,
15801592
const reply_callback_t &reply_callback);
15811593

1582-
client &xreadgroup(const std::vector<std::string> &cmd, const reply_callback_t &reply_callback);
1594+
client &xreadgroup(const std::string &group_name,
1595+
const std::string &consumer_name,
1596+
int count,
1597+
int block_milliseconds,
1598+
bool no_ack,
1599+
std::multimap<std::string, std::string> stream_members,
1600+
const reply_callback_t &reply_callback);
1601+
1602+
std::future<reply> xreadgroup(const std::string &group_name,
1603+
const std::string &consumer_name,
1604+
int count,
1605+
bool no_ack,
1606+
std::multimap<std::string, std::string> stream_members);
1607+
1608+
std::future<reply> xreadgroup(const std::string &group_name,
1609+
const std::string &consumer_name,
1610+
int count,
1611+
int block_milliseconds,
1612+
bool no_ack,
1613+
std::multimap<std::string, std::string> stream_members);
15831614

15841615
client &zadd(const std::string &key, const std::vector<std::string> &options,
15851616
const std::multimap<std::string, std::string> &score_members,

‎includes/cpp_redis/core/types.hpp‎

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,57 @@
1+
// The MIT License (MIT)
2+
//
3+
// Copyright (c) 2015-2017 Simon Ninon <simon.ninon@gmail.com>
4+
//
5+
// Permission is hereby granted, free of charge, to any person obtaining a copy
6+
// of this software and associated documentation files (the "Software"), to deal
7+
// in the Software without restriction, including without limitation the rights
8+
// to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
9+
// copies of the Software, and to permit persons to whom the Software is
10+
// furnished to do so, subject to the following conditions:
11+
//
12+
// The above copyright notice and this permission notice shall be included in all
13+
// copies or substantial portions of the Software.
14+
//
15+
// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
16+
// IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
17+
// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
18+
// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
19+
// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
20+
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
21+
// SOFTWARE.
22+
23+
#ifndef CPP_REDIS_TYPES_HPP
24+
#define CPP_REDIS_TYPES_HPP
25+
26+
#include <string>
27+
#include <vector>
28+
29+
30+
namespace cpp_redis {
31+
class range {
32+
public:
33+
enum class range_state {
34+
omit,
35+
include
36+
};
37+
explicit range(range_state state);
38+
explicit range(int count);
39+
range(int min, int max);
40+
range(int min, int max, int count);
41+
42+
bool should_omit() const;
43+
std::vector<std::string> get_args();
44+
std::vector<std::string> get_xpending_args() const;
45+
46+
private:
47+
int m_count;
48+
int m_min;
49+
int m_max;
50+
range_state m_state;
51+
};
52+
53+
typedef range range_t;
54+
} // namespace cpp_redis
55+
56+
57+
#endif //CPP_REDIS_TYPES_HPP

‎sources/core/client.cpp‎

Lines changed: 99 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -2222,6 +2222,7 @@ namespace cpp_redis {
22222222
return *this;
22232223
}
22242224

2225+
//<editor-fold desc="XGroup">
22252226
client &
22262227
client::xgroup_create(const std::string &key, const std::string &group_name, const reply_callback_t &reply_callback) {
22272228
send({"XGROUP", "CREATE", key, group_name, "$"}, reply_callback);
@@ -2261,6 +2262,7 @@ namespace cpp_redis {
22612262
send({"XGROUP", "DELCONSUMER", key, group_name, consumer_name}, reply_callback);
22622263
return *this;
22632264
}
2265+
//</editor-fold>
22642266

22652267
client &client::xinfo_consumers(const std::string &key, const std::string &group_name,
22662268
const reply_callback_t &reply_callback) {
@@ -2273,7 +2275,8 @@ namespace cpp_redis {
22732275
return *this;
22742276
}
22752277

2276-
client &client::xinfo_stream(const std::string &key, const reply_callback_t &reply_callback) {
2278+
client &
2279+
client::xinfo_stream(const std::string &key, const reply_callback_t &reply_callback) {
22772280
send({"XINFO", "STREAM", key}, reply_callback);
22782281
return *this;
22792282
}
@@ -2287,11 +2290,57 @@ namespace cpp_redis {
22872290
* @param reply_callback
22882291
* @return Integer reply: the number of entries of the stream at key.
22892292
*/
2290-
client &client::xlen(const std::string &key, const reply_callback_t &reply_callback) {
2293+
client &
2294+
client::xlen(const std::string &key, const reply_callback_t &reply_callback) {
22912295
send({"XLEN", key}, reply_callback);
22922296
return *this;
22932297
}
22942298

2299+
client &
2300+
client::xpending(const std::string &key,
2301+
const std::string &group_name,
2302+
const reply_callback_t &reply_callback) {
2303+
return xpending(key, group_name, range(range::range_state::omit), "", reply_callback);
2304+
}
2305+
2306+
client &
2307+
client::xpending(const std::string &key,
2308+
const std::string &group_name,
2309+
const range_t &range,
2310+
const reply_callback_t &reply_callback) {
2311+
return xpending(key, group_name, range, "", reply_callback);
2312+
}
2313+
2314+
client &
2315+
client::xpending(const std::string &key,
2316+
const std::string &group_name,
2317+
const std::string &consumer_name,
2318+
const reply_callback_t &reply_callback) {
2319+
return xpending(key, group_name, range(range::range_state::omit), consumer_name, reply_callback);
2320+
}
2321+
2322+
client &
2323+
client::xpending(const std::string &key,
2324+
const std::string &group_name,
2325+
const range_t &range,
2326+
const std::string &consumer_name,
2327+
const reply_callback_t &reply_callback) {
2328+
std::vector<std::string> cmd = {"XPENDING", key, group_name};
2329+
2330+
if (!range.should_omit()) {
2331+
auto range_members = range.get_xpending_args();
2332+
for (auto &rm : range_members) {
2333+
cmd.push_back(rm);
2334+
}
2335+
}
2336+
2337+
if (!consumer_name.empty()) {
2338+
cmd.push_back(consumer_name);
2339+
}
2340+
send(cmd, reply_callback);
2341+
return *this;
2342+
}
2343+
22952344
client &
22962345
client::xpending(const std::string &key, const std::string &group_name, const int &start, const int &end,
22972346
const int &count,
@@ -2316,17 +2365,41 @@ namespace cpp_redis {
23162365
client::xreadgroup(const std::string &group_name,
23172366
const std::string &consumer_name,
23182367
int count,
2319-
int block_milli_sec,
2368+
bool no_ack,
2369+
std::multimap<std::string, std::string> stream_members,
2370+
const reply_callback_t &reply_callback) {
2371+
std::vector<std::string> cmd = {"XREADGROUP",
2372+
"GROUP", group_name, consumer_name,
2373+
"COUNT", std::to_string(count)};
2374+
if (no_ack) {
2375+
cmd.push_back("NOACK");
2376+
}
2377+
2378+
//! score members
2379+
for (auto &sm : stream_members) {
2380+
cmd.push_back(sm.first);
2381+
cmd.push_back(sm.second);
2382+
}
2383+
2384+
send(cmd, reply_callback);
2385+
return *this;
2386+
}
2387+
2388+
client &
2389+
client::xreadgroup(const std::string &group_name,
2390+
const std::string &consumer_name,
2391+
int count,
2392+
int block_milliseconds,
23202393
bool no_ack,
23212394
std::multimap<std::string, std::string> stream_members,
23222395
const reply_callback_t &reply_callback) {
23232396
std::vector<std::string> cmd = {"XREADGROUP",
23242397
"GROUP", group_name, consumer_name,
23252398
"COUNT", std::to_string(count)};
23262399

2327-
if (block_milli_sec >= 0) {
2400+
if (block_milliseconds >= 0) {
23282401
cmd.emplace_back("BLOCK");
2329-
cmd.push_back(std::to_string(block_milli_sec));
2402+
cmd.push_back(std::to_string(block_milliseconds));
23302403
}
23312404

23322405
if (no_ack) {
@@ -4225,6 +4298,27 @@ namespace cpp_redis {
42254298
});
42264299
}
42274300

4301+
std::future<reply> client::xreadgroup(const std::string &group_name,
4302+
const std::string &consumer_name,
4303+
int count,
4304+
bool no_ack,
4305+
std::multimap<std::string, std::string> stream_members) {
4306+
return exec_cmd([=](const reply_callback_t &cb) -> client & {
4307+
return xreadgroup(group_name, consumer_name, count, no_ack, stream_members, cb);
4308+
});
4309+
}
4310+
4311+
std::future<reply> client::xreadgroup(const std::string &group_name,
4312+
const std::string &consumer_name,
4313+
int count,
4314+
int block_milliseconds,
4315+
bool no_ack,
4316+
std::multimap<std::string, std::string> stream_members) {
4317+
return exec_cmd([=](const reply_callback_t &cb) -> client & {
4318+
return xreadgroup(group_name, consumer_name, count, block_milliseconds, no_ack, stream_members, cb);
4319+
});
4320+
}
4321+
42284322
std::future<reply>
42294323
client::zadd(const std::string &key, const std::vector<std::string> &options,
42304324
const std::multimap<std::string, std::string> &score_members) {
@@ -4580,17 +4674,4 @@ namespace cpp_redis {
45804674
});
45814675
}
45824676

4583-
client &client::xreadgroup(const std::vector<std::string> &cmd, const reply_callback_t &reply_callback) {
4584-
/*std::vector<std::string> cmd = {"XADD", key, id};
4585-
4586-
//! score members
4587-
for (auto &sm : field_members) {
4588-
cmd.push_back(sm.first);
4589-
cmd.push_back(sm.second);
4590-
}*/
4591-
4592-
send(cmd, reply_callback);
4593-
return *this;
4594-
}
4595-
45964677
} // namespace cpp_redis

‎sources/core/types.cpp‎

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
/*
2+
*
3+
* Created by nick on 11/23/18.
4+
*
5+
* Copyright(c) 2018 Iris. All rights reserved.
6+
*
7+
* Use and copying of this software and preparation of derivative
8+
* works based upon this software are not permitted. Any distribution
9+
* of this software or derivative works must comply with all applicable
10+
* Canadian export control laws.
11+
*
12+
* THIS SOFTWARE IS PROVIDED ``AS IS'' AND ANY EXPRESSED OR IMPLIED
13+
* WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES
14+
* OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
15+
* DISCLAIMED. IN NO EVENT SHALL IRIS OR ITS CONTRIBUTORS
16+
* BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY,
17+
* OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
18+
* PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA,
19+
* OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON
20+
* ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY,
21+
* OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT
22+
* OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF
23+
* SUCH DAMAGE.
24+
*
25+
*/
26+
27+
28+
#include <cpp_redis/core/types.hpp>
29+
30+
namespace cpp_redis {
31+
range::range() = default;
32+
33+
range::range(int count) : m_min(-1) {}
34+
35+
range::range(range_state state) : m_state(state) {}
36+
37+
range::range(int min, int max) : m_count(10) {}
38+
39+
range::range(int min, int max, int count) {}
40+
41+
bool range::should_omit() const {
42+
return m_state == range_state::omit;
43+
}
44+
45+
std::vector<std::string> range::get_args() {}
46+
47+
std::vector<std::string> range::get_xpending_args() const {
48+
return m_min == -1
49+
? std::vector<std::string>{"-",
50+
"+",
51+
std::to_string(m_count)}
52+
: std::vector<std::string>{std::to_string(m_min),
53+
std::to_string(m_max),
54+
std::to_string(m_count)};
55+
}
56+
} // namespace cpp_redis

0 commit comments

Comments
 (0)

Back | FazBrowse Home | New Git URL