| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -32,6 +32,7 @@ | |||
| 32 | 32 | #include <string> | |
| 33 | 33 | #include <vector> | |
| 34 | 34 | ||
| 35 | + #include <cpp_redis/core/types.hpp> | ||
| 35 | 36 | #include <cpp_redis/core/sentinel.hpp> | |
| 36 | 37 | #include <cpp_redis/helpers/variadic_template.hpp> | |
| 37 | 38 | #include <cpp_redis/misc/logger.hpp> | |
@@ -1563,23 +1564,53 @@ namespace cpp_redis { | |||
| 1563 | 1564 | ||
| 1564 | 1565 | client &xpending(const std::string &key, | |
| 1565 | 1566 | 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, | ||
| 1570 | 1583 | const reply_callback_t &reply_callback = nullptr); | |
| 1571 | 1584 | ||
| 1572 | 1585 | //XREADGROUP GROUP group consumer [COUNT count] [BLOCK milliseconds] [NOACK] STREAMS key [key ...] ID [ID ...] | |
| 1573 | 1586 | ||
| 1574 | 1587 | client &xreadgroup(const std::string &group_name, | |
| 1575 | 1588 | const std::string &consumer_name, | |
| 1576 | 1589 | int count, | |
| 1577 | - int block_milli_sec, | ||
| 1578 | 1590 | bool no_ack, | |
| 1579 | 1591 | std::multimap<std::string, std::string> stream_members, | |
| 1580 | 1592 | const reply_callback_t &reply_callback); | |
| 1581 | 1593 | ||
| 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); | ||
| 1583 | 1614 | ||
| 1584 | 1615 | client &zadd(const std::string &key, const std::vector<std::string> &options, | |
| 1585 | 1616 | const std::multimap<std::string, std::string> &score_members, | |
| Original file line number | Diff line number | Diff 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 | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -2222,6 +2222,7 @@ namespace cpp_redis { | |||
| 2222 | 2222 | return *this; | |
| 2223 | 2223 | } | |
| 2224 | 2224 | ||
| 2225 | + //<editor-fold desc="XGroup"> | ||
| 2225 | 2226 | client & | |
| 2226 | 2227 | client::xgroup_create(const std::string &key, const std::string &group_name, const reply_callback_t &reply_callback) { | |
| 2227 | 2228 | send({"XGROUP", "CREATE", key, group_name, "$"}, reply_callback); | |
@@ -2261,6 +2262,7 @@ namespace cpp_redis { | |||
| 2261 | 2262 | send({"XGROUP", "DELCONSUMER", key, group_name, consumer_name}, reply_callback); | |
| 2262 | 2263 | return *this; | |
| 2263 | 2264 | } | |
| 2265 | + //</editor-fold> | ||
| 2264 | 2266 | ||
| 2265 | 2267 | client &client::xinfo_consumers(const std::string &key, const std::string &group_name, | |
| 2266 | 2268 | const reply_callback_t &reply_callback) { | |
@@ -2273,7 +2275,8 @@ namespace cpp_redis { | |||
| 2273 | 2275 | return *this; | |
| 2274 | 2276 | } | |
| 2275 | 2277 | ||
| 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) { | ||
| 2277 | 2280 | send({"XINFO", "STREAM", key}, reply_callback); | |
| 2278 | 2281 | return *this; | |
| 2279 | 2282 | } | |
@@ -2287,11 +2290,57 @@ namespace cpp_redis { | |||
| 2287 | 2290 | * @param reply_callback | |
| 2288 | 2291 | * @return Integer reply: the number of entries of the stream at key. | |
| 2289 | 2292 | */ | |
| 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) { | ||
| 2291 | 2295 | send({"XLEN", key}, reply_callback); | |
| 2292 | 2296 | return *this; | |
| 2293 | 2297 | } | |
| 2294 | 2298 | ||
| 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 | + | ||
| 2295 | 2344 | client & | |
| 2296 | 2345 | client::xpending(const std::string &key, const std::string &group_name, const int &start, const int &end, | |
| 2297 | 2346 | const int &count, | |
@@ -2316,17 +2365,41 @@ namespace cpp_redis { | |||
| 2316 | 2365 | client::xreadgroup(const std::string &group_name, | |
| 2317 | 2366 | const std::string &consumer_name, | |
| 2318 | 2367 | 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, | ||
| 2320 | 2393 | bool no_ack, | |
| 2321 | 2394 | std::multimap<std::string, std::string> stream_members, | |
| 2322 | 2395 | const reply_callback_t &reply_callback) { | |
| 2323 | 2396 | std::vector<std::string> cmd = {"XREADGROUP", | |
| 2324 | 2397 | "GROUP", group_name, consumer_name, | |
| 2325 | 2398 | "COUNT", std::to_string(count)}; | |
| 2326 | 2399 | ||
| 2327 | - if (block_milli_sec >= 0) { | ||
| 2400 | + if (block_milliseconds >= 0) { | ||
| 2328 | 2401 | cmd.emplace_back("BLOCK"); | |
| 2329 | - cmd.push_back(std::to_string(block_milli_sec)); | ||
| 2402 | + cmd.push_back(std::to_string(block_milliseconds)); | ||
| 2330 | 2403 | } | |
| 2331 | 2404 | ||
| 2332 | 2405 | if (no_ack) { | |
@@ -4225,6 +4298,27 @@ namespace cpp_redis { | |||
| 4225 | 4298 | }); | |
| 4226 | 4299 | } | |
| 4227 | 4300 | ||
| 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 | + | ||
| 4228 | 4322 | std::future<reply> | |
| 4229 | 4323 | client::zadd(const std::string &key, const std::vector<std::string> &options, | |
| 4230 | 4324 | const std::multimap<std::string, std::string> &score_members) { | |
@@ -4580,17 +4674,4 @@ namespace cpp_redis { | |||
| 4580 | 4674 | }); | |
| 4581 | 4675 | } | |
| 4582 | 4676 | ||
| 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 | - | ||
| 4596 | 4677 | } // namespace cpp_redis | |
| Original file line number | Diff line number | Diff 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 | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments