FazBrowse GitHub Viewer
|
Trending
|
URL:
|
Home
Tools:
[Download Repo ZIP]
[View Raw Code]
[Original HTTPS Page]
cpp_redis/examples/cpp_redis_streams_client.cpp at master · cpp-redis/cpp_redis · GitHub
cpp-redis
cpp_redis
Repository navigation
Code
Issues
35
(35)
Pull requests
9
(9)
Discussions
Actions
Projects
Wiki
Security and quality
Insights
Expand file tree
Breadcrumbs
cpp_redis
/
examples
/
cpp_redis_streams_client.cpp
Copy path
More file actions
More file actions
Latest commit
History
History
History
119 lines (95 loc) · 4.23 KB
Breadcrumbs
cpp_redis
/
examples
/
cpp_redis_streams_client.cpp
Copy path
File metadata and controls
119 lines (95 loc) · 4.23 KB
Raw
Copy raw file
Download raw file
Open symbols panel
Edit and raw actions
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
//
The MIT License (MIT)
//
//
Copyright (c) 2015-2017 Simon Ninon <simon.ninon@gmail.com>
//
//
Permission is hereby granted, free of charge, to any person obtaining a copy
//
of this software and associated documentation files (the "Software"), to deal
//
in the Software without restriction, including without limitation the rights
//
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
//
copies of the Software, and to permit persons to whom the Software is
//
furnished to do so, subject to the following conditions:
//
//
The above copyright notice and this permission notice shall be included in all
//
copies or substantial portions of the Software.
//
//
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
//
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
//
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
//
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
//
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
//
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
//
SOFTWARE.
#
include
<
string
>
#
include
<
cpp_redis/cpp_redis
>
#
include
"
winsock_initializer.h
"
int
main
() {
winsock_initializer winsock_init;
//
! Enable logging
cpp_redis::active_logger = std::unique_ptr<cpp_redis::logger>(
new
cpp_redis::logger);
cpp_redis::client client;
client.
connect
(
"
127.0.0.1
"
,
6379
,
[](
const
std::string &host, std::
size_t
port, cpp_redis::connect_state status) {
if
(status == cpp_redis::connect_state::dropped) {
std::cout <<
"
client disconnected from
"
<< host <<
"
:
"
<< port << std::endl;
}
});
auto
reply_cmd = [](cpp_redis::reply &reply) {
std::cout <<
"
response:
"
<< reply.
as_string
() << std::endl;
};
std::string message_id;
const
std::string group_name =
"
groupone
"
;
const
std::string session_name =
"
sessone
"
;
const
std::string consumer_name =
"
ABCD
"
;
std::multimap<std::string, std::string> ins;
ins.
insert
(std::pair<std::string, std::string>{
"
message
"
,
"
hello
"
});
ins.
insert
(std::pair<std::string, std::string>{
"
result
"
,
"
a result
"
});
client.
xtrim
(session_name,
10
, reply_cmd);
client.
xgroup_create
(session_name, group_name, reply_cmd);
client.
xadd
(session_name,
"
*
"
, {{
"
message
"
,
"
hello
"
},
{
"
details
"
,
"
some details
"
}}, [&](cpp_redis::reply &reply) {
std::cout <<
"
response:
"
<< reply.
as_string
() << std::endl;
message_id = reply.
as_string
();
std::cout <<
"
message id:
"
<< message_id << std::endl;
});
client.
sync_commit
();
std::cout <<
"
message id after:
"
<< message_id << std::endl;
client.
xack
(session_name, group_name, {message_id}, reply_cmd);
client.
xinfo_stream
(session_name, [](cpp_redis::reply &reply) {
//
std::cout << reply << std::endl;
cpp_redis::xinfo_reply
x
(reply);
std::cout <<
"
Len:
"
<< x.
Length
<< std::endl;
});
//
client.xadd(session_name, message_id, {{"final", "finished"}}, reply_cmd);
client.
sync_commit
(
std::chrono::milliseconds
(
100
));
client.
xread
({Streams: {{session_name},
{message_id}},
Count:
10
,
Block:
100
}, reply_cmd);
client.
sync_commit
(
std::chrono::milliseconds
(
100
));
client.
xrange
(session_name, {
"
-
"
,
"
+
"
,
10
}, reply_cmd);
client.
xreadgroup
({group_name,
"
0
"
,
{{session_name}, {
"
>
"
}},
1
, -
1
,
false
//
count, block, no_ack
}, [](cpp_redis::reply &reply) {
cpp_redis::xstream_reply
msg
(reply);
std::cout << msg << std::endl;
});
client.
sync_commit
(
std::chrono::milliseconds
(
100
));
client.
xreadgroup
({group_name,
consumer_name,
{{session_name}, {
"
>
"
}},
1
,
0
,
false
//
count, block, no_ack
}, [](cpp_redis::reply &reply) {
cpp_redis::xstream_reply
msg
(reply);
std::cout << msg << std::endl;
});
//
commands are pipelined and only sent when client.commit() is called
//
client.commit();
//
synchronous commit, no timeout
client.
sync_commit
(
std::chrono::milliseconds
(
100
));
//
synchronous commit, timeout
//
client.sync_commit(std::chrono::milliseconds(100));
return
0
;
}
Back
|
FazBrowse Home
|
New Git URL