Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
C
cmu
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
李维杰
cmu
Commits
40cf7643
Commit
40cf7643
authored
Jun 22, 2020
by
李维杰
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
重构请求相关逻辑
parent
f134760d
Hide whitespace changes
Inline
Side-by-side
Showing
6 changed files
with
389 additions
and
136 deletions
+389
-136
cmd_processer_server.cc
cmu/src/cmd_processer_server.cc
+62
-33
cmd_processer_server.h
cmu/src/cmd_processer_server.h
+4
-1
cmd_processer_client.cc
kcp-client/src/cmd_processer_client.cc
+100
-11
cmd_processer_client.h
kcp-client/src/cmd_processer_client.h
+16
-1
mqtt_request.cc
thirdparty/mqtt/mqtt_request.cc
+137
-83
mqtt_request.h
thirdparty/mqtt/mqtt_request.h
+70
-7
No files found.
cmu/src/cmd_processer_server.cc
View file @
40cf7643
...
...
@@ -54,63 +54,93 @@ namespace offcn
{
printf
(
"
\n
receive message : topic = %s, context : %s
\n
"
,
topic
.
c_str
(),
message
.
c_str
());
std
::
string
path
;
std
::
string
method
;
std
::
string
error
;
ReqValue
values
;
if
(
!
MqttRequest
::
ParseRequest
(
message
,
path
,
values
))
if
(
!
MqttRequest
::
ParseRequest
Type
(
message
,
method
,
error
))
{
printf
(
"Error : Parse request error
\n
"
);
return
;
}
std
::
string
fid
=
values
[
"fid"
]
;
std
::
string
pull_type
;
if
(
path
.
compare
(
"join"
)
==
0
)
ReqValue
valuse
;
if
(
method
.
compare
(
"join"
)
==
0
)
{
pull_type
=
values
[
"pull_type"
];
OnRequstJoin
(
fid
,
pull_type
);
if
(
MqttRequest
::
ParseJoinRequest
(
message
,
valuse
,
error
))
{
OnParseJoinSuccess
(
values
[
"message_id"
],
valuse
[
"from_id"
]);
}
else
{
OnParseFailure
(
values
[
"message_id"
],
valuse
[
"from_id"
],
error
);
}
}
else
if
(
method
.
compare
(
"update"
)
==
0
)
{
if
(
MqttRequest
::
ParseUpdateRequest
(
message
,
valuse
,
error
))
{
OnParseUpdateSuccess
(
valuse
[
"message_id"
],
values
[
"from_id"
],
values
[
"pull_type"
]);
}
else
{
OnParseFailure
(
valuse
[
"message_id"
],
values
[
"from_id"
],
error
);
}
}
}
void
CmdProcesserServer
::
On
RequstJoin
(
std
::
string
fid
,
std
::
string
pull_type
)
void
CmdProcesserServer
::
On
ParseJoinSuccess
(
std
::
string
message_id
,
std
::
string
from_id
)
{
if
(
client_lists_
.
find
(
fid
)
!=
client_lists_
.
end
())
std
::
string
response
;
if
(
client_lists_
.
find
(
from_id
)
!=
client_lists_
.
end
())
{
printf
(
"Waring : fid = %s already exist
\n
"
,
fid
.
c_str
());
response
=
MqttRequest
::
BuildResponse
(
message_id
,
"id exist"
);
mqtt_wrapper_
->
SendRequest
(
from_id
,
response
);
return
;
}
std
::
string
dest_id
;
std
::
map
<
std
::
string
,
std
::
string
>::
iterator
itor
;
for
(
itor
=
client_lists_
.
begin
();
itor
!=
client_lists_
.
end
();
itor
++
)
else
{
if
(
itor
->
second
.
compare
(
"rtmp"
)
==
0
)
client_lists_
.
insert
(
std
::
pair
<
std
::
string
,
std
::
string
>
(
from_id
,
"null"
));
response
=
MqttRequest
::
BuildResponse
(
message_id
,
"ok"
);
printf
(
"
\n
------------------------------------
\n
"
);
std
::
map
<
std
::
string
,
std
::
string
>::
iterator
itor
;
for
(
itor
=
client_lists_
.
begin
();
itor
!=
client_lists_
.
end
();
itor
++
)
{
dest_id
=
itor
->
first
;
break
;
printf
(
"id = %s, pull type = %s
\n
"
,
itor
->
first
.
c_str
(),
itor
->
second
.
c_str
());
}
printf
(
"------------------------------------
\n
"
);
}
}
void
CmdProcesserServer
::
OnParseFailure
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
error
)
{
std
::
string
response
=
MqttRequest
::
BuildResponse
(
message_id
,
error
);
mqtt_wrapper_
->
SendRequest
(
from_id
,
response
);
}
void
CmdProcesserServer
::
OnParseUpdateSuccess
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
pull_type
)
{
std
::
string
response
;
if
(
itor
!=
client_lists_
.
end
())
if
(
client_lists_
.
find
(
from_id
)
==
client_lists_
.
end
())
{
client_lists_
.
insert
(
std
::
pair
<
std
::
string
,
std
::
string
>
(
fid
,
"kcp"
));
response
=
MqttRequest
::
JoinResponse
(
dest_id
,
"success"
);
response
=
MqttRequest
::
BuildResponse
(
message_id
,
"id is not exist"
);
mqtt_wrapper_
->
SendRequest
(
from_id
,
response
);
return
;
}
else
{
client_lists_
.
insert
(
std
::
pair
<
std
::
string
,
std
::
string
>
(
fid
,
"rtmp"
));
response
=
MqttRequest
::
JoinResponse
(
""
,
"success"
);
}
mqtt_wrapper_
->
SendRequest
(
fid
,
response
);
client_lists_
.
insert
(
std
::
pair
<
std
::
string
,
std
::
string
>
(
from_id
,
pull_type
));
response
=
MqttRequest
::
BuildResponse
(
message_id
,
"ok"
);
printf
(
"
\n
------------------------------------
\n
"
);
for
(
itor
=
client_lists_
.
begin
();
itor
!=
client_lists_
.
end
();
itor
++
)
{
printf
(
"id = %s, pull type = %s
\n
"
,
itor
->
first
.
c_str
(),
itor
->
second
.
c_str
());
printf
(
"
\n
------------------------------------
\n
"
);
std
::
map
<
std
::
string
,
std
::
string
>::
iterator
itor
;
for
(
itor
=
client_lists_
.
begin
();
itor
!=
client_lists_
.
end
();
itor
++
)
{
printf
(
"id = %s, pull type = %s
\n
"
,
itor
->
first
.
c_str
(),
itor
->
second
.
c_str
());
}
printf
(
"------------------------------------
\n
"
);
}
printf
(
"------------------------------------
\n
"
);
}
}
\ No newline at end of file
cmu/src/cmd_processer_server.h
View file @
40cf7643
...
...
@@ -26,7 +26,10 @@ namespace offcn
virtual
void
OnServerMessageArrived
(
std
::
string
topic
,
std
::
string
message
);
private
:
void
OnRequstJoin
(
std
::
string
fid
,
std
::
string
pull_type
);
void
OnParseJoinSuccess
(
std
::
string
message_id
,
std
::
string
from_id
);
void
OnParseFailure
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
error
);
void
OnParseUpdateSuccess
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
pull_type
);
private
:
MqttWrapper
*
mqtt_wrapper_
;
...
...
kcp-client/src/cmd_processer_client.cc
View file @
40cf7643
#include "cmd_processer_client.h"
#include "mqtt_request.h"
#include <time.h>
namespace
offcn
{
static
const
std
::
string
kGlobalTopic
=
"kcp_server_mqtt"
;
static
const
std
::
string
kLocalID
=
"mqtt_client"
;
static
const
std
::
string
kDstTopic
=
kLocalID
;
static
std
::
string
clientMessageID
()
{
char
str
[
12
]
=
{
0
};
int
len
=
12
;
srand
((
unsigned
int
)
time
(
NULL
));
int
i
;
for
(
i
=
0
;
i
<
len
;
++
i
)
{
switch
((
rand
()
%
3
))
{
case
1
:
str
[
i
]
=
'A'
+
rand
()
%
26
;
break
;
case
2
:
str
[
i
]
=
'a'
+
rand
()
%
26
;
break
;
default
:
str
[
i
]
=
'0'
+
rand
()
%
10
;
break
;
}
}
str
[
++
i
]
=
'\0'
;
return
str
;
}
CmdProcesserClient
::
CmdProcesserClient
()
:
mqtt_wrapper_
(
NULL
)
:
mqtt_wrapper_
(
NULL
),
observer_
(
NULL
)
{
local_id_
=
kLocalID
;
...
...
@@ -23,9 +52,10 @@ namespace offcn
}
}
void
CmdProcesserClient
::
SetLocalID
(
std
::
string
id
)
void
CmdProcesserClient
::
SetLocalID
(
std
::
string
id
,
CmdProcesserObserver
*
observer
)
{
local_id_
=
id
;
observer_
=
observer
;
}
bool
CmdProcesserClient
::
Connect
()
{
...
...
@@ -37,7 +67,16 @@ namespace offcn
}
bool
CmdProcesserClient
::
SendLocalCandidate
(
std
::
string
tid
,
std
::
string
type
,
std
::
string
candidate
)
{
std
::
string
req
=
MqttRequest
::
Candidate
(
local_id_
,
tid
,
type
,
candidate
);
std
::
string
req
;
if
(
type
.
compare
(
"offer"
)
==
0
)
{
req
=
MqttRequest
::
OfferRequest
(
clientMessageID
(),
local_id_
,
tid
,
candidate
);
}
else
if
(
type
.
compare
(
"answer"
)
==
0
)
{
req
=
MqttRequest
::
AnswerRequest
(
clientMessageID
(),
local_id_
,
tid
,
candidate
);
}
return
mqtt_wrapper_
->
SendRequest
(
tid
,
req
);
}
...
...
@@ -58,7 +97,7 @@ namespace offcn
printf
(
"subscribe topic success!
\n
"
);
std
::
string
topic
=
kGlobalTopic
;
std
::
string
req
=
MqttRequest
::
JoinRequest
(
local_id_
);
;
std
::
string
req
=
MqttRequest
::
JoinRequest
(
clientMessageID
(),
local_id_
)
;
mqtt_wrapper_
->
SendRequest
(
topic
,
req
);
}
...
...
@@ -68,22 +107,71 @@ namespace offcn
}
void
CmdProcesserClient
::
OnServerMessageArrived
(
std
::
string
topic
,
std
::
string
message
)
{
std
::
string
path
;
std
::
string
method
;
std
::
string
error
;
ReqValue
values
;
if
(
!
MqttRequest
::
ParseRe
sponse
(
message
,
path
,
values
))
if
(
!
MqttRequest
::
ParseRe
questType
(
message
,
method
,
error
))
{
printf
(
"Error : Parse request error
\n
"
);
return
;
}
if
(
path
.
compare
(
"join"
)
==
0
)
if
(
method
.
compare
(
"publishers"
)
==
0
)
{
if
(
MqttRequest
::
ParsePublishersRequest
(
message
,
values
,
error
))
{
OnParsePublishersSuccess
(
values
[
"message_id"
],
values
[
"publishers_id"
]);
}
else
{
OnParseFailure
(
values
[
"message_id"
],
values
[
"publishers_id"
],
error
);
}
}
else
if
(
method
.
compare
(
"offer"
)
==
0
)
{
std
::
string
ids
=
values
[
"ids"
];
std
::
string
result
=
values
[
"result"
];
if
(
ids
.
length
()
>
0
)
if
(
MqttRequest
::
ParseOfferRequest
(
message
,
values
,
error
))
{
publisher_list_
.
insert
(
std
::
pair
<
std
::
string
,
std
::
string
>
(
ids
,
ids
)
);
OnParseOfferOrAnswerSuccess
(
values
[
"message_id"
],
"offer"
,
values
[
"from_id"
],
values
[
"candidate"
]
);
}
else
{
OnParseFailure
(
values
[
"message_id"
],
values
[
"from_id"
],
error
);
}
}
else
if
(
method
.
compare
(
"answer"
)
==
0
)
{
if
(
MqttRequest
::
ParseAnswerRequest
(
message
,
values
,
error
))
{
OnParseOfferOrAnswerSuccess
(
values
[
"message_id"
],
"answer"
,
values
[
"from_id"
],
values
[
"candidate"
]);
}
else
{
OnParseFailure
(
values
[
"message_id"
],
values
[
"from_id"
],
error
);
}
}
}
void
CmdProcesserClient
::
OnParsePublishersSuccess
(
std
::
string
message_id
,
std
::
string
publishers_id
)
{
publisher_list_
.
insert
(
std
::
pair
<
std
::
string
,
std
::
string
>
(
publishers_id
,
publishers_id
));
std
::
string
req
=
MqttRequest
::
BuildResponse
(
message_id
,
"ok"
);
std
::
string
topic
=
kGlobalTopic
;
mqtt_wrapper_
->
SendRequest
(
topic
,
req
);
}
void
CmdProcesserClient
::
OnParseFailure
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
error
)
{
std
::
string
response
=
MqttRequest
::
BuildResponse
(
message_id
,
error
);
mqtt_wrapper_
->
SendRequest
(
from_id
,
response
);
}
void
CmdProcesserClient
::
OnParseOfferOrAnswerSuccess
(
std
::
string
message_id
,
std
::
string
type
,
std
::
string
to_id
,
std
::
string
candidate
)
{
std
::
string
req
=
MqttRequest
::
BuildResponse
(
message_id
,
"ok"
);
mqtt_wrapper_
->
SendRequest
(
to_id
,
req
);
if
(
observer_
)
{
observer_
->
OnReceiveRemoteCandidate
(
to_id
,
type
,
candidate
);
}
}
}
\ No newline at end of file
kcp-client/src/cmd_processer_client.h
View file @
40cf7643
...
...
@@ -5,6 +5,12 @@
namespace
offcn
{
class
CmdProcesserObserver
{
public
:
virtual
void
OnReceiveRemoteCandidate
(
std
::
string
&
remote_id
,
std
::
string
&
type
,
std
::
string
&
candidate
)
=
0
;
};
class
CmdProcesserClient
:
public
CmdObserver
{
public
:
...
...
@@ -12,7 +18,7 @@ namespace offcn
~
CmdProcesserClient
();
public
:
void
SetLocalID
(
std
::
string
id
);
void
SetLocalID
(
std
::
string
id
,
CmdProcesserObserver
*
observer
);
bool
Connect
();
bool
SubscribeTopic
();
/**
...
...
@@ -30,6 +36,12 @@ namespace offcn
virtual
void
OnServerSubscribeFailure
();
virtual
void
OnServerMessageArrived
(
std
::
string
topic
,
std
::
string
message
);
public
:
void
OnParsePublishersSuccess
(
std
::
string
message_id
,
std
::
string
publishers_id
);
void
OnParseFailure
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
error
);
void
OnParseOfferOrAnswerSuccess
(
std
::
string
message_id
,
std
::
string
type
,
std
::
string
to_id
,
std
::
string
candidate
);
private
:
MqttWrapper
*
mqtt_wrapper_
;
private
:
...
...
@@ -37,5 +49,7 @@ namespace offcn
std
::
map
<
std
::
string
,
std
::
string
>
subscriber_list_
;
private
:
std
::
string
local_id_
;
private
:
CmdProcesserObserver
*
observer_
;
};
}
\ No newline at end of file
thirdparty/mqtt/mqtt_request.cc
View file @
40cf7643
...
...
@@ -3,154 +3,208 @@
namespace
offcn
{
/*
static
bool
GetJsonValue
(
std
::
string
&
response
,
Json
::
Value
&
root
)
{
"path": "join",
"fid" : "xxxid",
"data" :
{
"pull_type" : "null"
}
Json
::
CharReaderBuilder
readerBuilder
;
std
::
unique_ptr
<
Json
::
CharReader
>
const
reader
(
readerBuilder
.
newCharReader
());
const
char
*
start_pos
=
response
.
c_str
();
std
::
string
err
;
if
(
!
reader
->
parse
(
start_pos
,
start_pos
+
response
.
length
(),
&
root
,
&
err
))
{
return
false
;
}
return
true
;
}
*/
std
::
string
MqttRequest
::
JoinRequest
(
std
::
string
local
_id
)
std
::
string
MqttRequest
::
JoinRequest
(
std
::
string
message_id
,
std
::
string
from
_id
)
{
Json
::
Value
root
;
Json
::
StreamWriterBuilder
wBuilder
;
root
[
"path"
]
=
"join"
;
root
[
"fid"
]
=
local_id
;
Json
::
Value
data
;
data
[
"pull_type"
]
=
"null"
;
root
[
"data"
]
=
data
;
root
[
"message_id"
]
=
message_id
;
root
[
"method"
]
=
"join"
;
root
[
"from_id"
]
=
from_id
;
return
Json
::
writeString
(
wBuilder
,
root
);
}
/*
{
"path": "update",
"fid" : "xxxid",
"data" :
{
"pull_type" : "rtmp" or "null" or "kcp"
}
}
*/
std
::
string
MqttRequest
::
UpdateRequest
(
std
::
string
local_id
,
std
::
string
pull_type
)
std
::
string
MqttRequest
::
UpdateRequest
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
pull_type
)
{
Json
::
Value
root
;
Json
::
StreamWriterBuilder
wBuilder
;
root
[
"path"
]
=
"update"
;
root
[
"fid"
]
=
local_id
;
root
[
"message_id"
]
=
message_id
;
root
[
"method"
]
=
"update"
;
root
[
"from_id"
]
=
from_id
;
root
[
"pull_type"
]
=
pull_type
;
Json
::
Value
data
;
data
[
"pull_type"
]
=
pull_type
;
return
Json
::
writeString
(
wBuilder
,
root
);
}
std
::
string
MqttRequest
::
OfferRequest
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
to_id
,
std
::
string
candidate
)
{
Json
::
Value
root
;
Json
::
StreamWriterBuilder
wBuilder
;
root
[
"data"
]
=
data
;
root
[
"message_id"
]
=
message_id
;
root
[
"method"
]
=
"offer"
;
root
[
"from_id"
]
=
from_id
;
root
[
"to_id"
]
=
to_id
;
root
[
"candidate"
]
=
candidate
;
return
Json
::
writeString
(
wBuilder
,
root
);
}
/*
std
::
string
MqttRequest
::
AnswerRequest
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
to_id
,
std
::
string
candidate
)
{
"path": "offer/answer",
"fid" : "xxxid",
"tid" : "remoteid",
"candidate" : ""
Json
::
Value
root
;
Json
::
StreamWriterBuilder
wBuilder
;
root
[
"message_id"
]
=
message_id
;
root
[
"method"
]
=
"answer"
;
root
[
"from_id"
]
=
from_id
;
root
[
"to_id"
]
=
to_id
;
root
[
"candidate"
]
=
candidate
;
return
Json
::
writeString
(
wBuilder
,
root
);
}
*/
std
::
string
MqttRequest
::
Candidate
(
std
::
string
local_id
,
std
::
string
remote_id
,
std
::
string
type
,
std
::
string
candidate
)
std
::
string
MqttRequest
::
PublishersRequest
(
std
::
string
message_id
,
std
::
string
to_id
,
std
::
string
publisher_id
)
{
Json
::
Value
root
;
Json
::
StreamWriterBuilder
wBuilder
;
root
[
"
path"
]
=
type
;
root
[
"
fid"
]
=
local_id
;
root
[
"
tid"
]
=
remote
_id
;
root
[
"
candidate"
]
=
candidate
;
root
[
"
message_id"
]
=
message_id
;
root
[
"
method"
]
=
"publishers"
;
root
[
"
publishers_id"
]
=
publisher
_id
;
root
[
"
to_id"
]
=
to_id
;
return
Json
::
writeString
(
wBuilder
,
root
);
}
static
bool
GetJsonValue
(
std
::
string
&
response
,
Json
::
Value
&
root
)
bool
MqttRequest
::
ParseRequestType
(
std
::
string
request
,
std
::
string
&
type
,
std
::
string
&
error
)
{
Json
::
CharReaderBuilder
readerBuilder
;
std
::
unique_ptr
<
Json
::
CharReader
>
const
reader
(
readerBuilder
.
newCharReader
());
const
char
*
start_pos
=
response
.
c_str
();
std
::
string
err
;
if
(
!
reader
->
parse
(
start_pos
,
start_pos
+
response
.
length
(),
&
root
,
&
err
))
Json
::
Value
root
;
if
(
!
GetJsonValue
(
request
,
root
))
{
error
=
"json format error"
;
return
false
;
}
if
(
!
root
[
"method"
])
{
error
=
"param error"
;
return
false
;
}
return
true
;
}
bool
MqttRequest
::
Parse
Request
(
std
::
string
&
req
,
std
::
string
&
path
,
ReqValue
&
values
)
bool
MqttRequest
::
Parse
JoinRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
)
{
Json
::
Value
root
;
if
(
!
GetJsonValue
(
request
,
root
))
{
error
=
"json format error"
;
goto
Error
;
}
if
(
!
root
[
"message_id"
]
||
!
root
[
"method"
]
||
!
root
[
"from_id"
])
{
error
=
"param error"
;
goto
Error
;
}
values
[
"message_id"
]
=
root
[
"message_id"
].
asString
();
values
[
"method"
]
=
root
[
"method"
].
asString
();
values
[
"from_id"
]
=
root
[
"from_id"
].
asString
();
if
(
!
GetJsonValue
(
req
,
root
))
return
false
;
if
(
!
root
[
"path"
]
||
!
root
[
"fid"
])
return
false
;
return
true
;
path
=
root
[
"path"
].
asString
();
std
::
string
fid
=
root
[
"fid"
].
asString
();
Error
:
values
[
"message_id"
]
=
""
;
values
[
"method"
]
=
""
;
values
[
"from_id"
]
=
""
;
if
(
path
.
compare
(
"join"
)
==
0
||
path
.
compare
(
"update"
)
==
0
)
return
false
;
}
bool
MqttRequest
::
ParseUpdateRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
)
{
Json
::
Value
root
;
if
(
!
GetJsonValue
(
request
,
root
))
{
values
[
"fid"
]
=
fid
;
values
[
"pull_type"
]
=
root
[
"pull_type"
].
asString
()
;
error
=
"json format error"
;
return
false
;
}
else
if
(
path
.
compare
(
"offer"
)
==
0
||
path
.
compare
(
"answer"
)
==
0
)
if
(
!
root
[
"message_id"
]
||
!
root
[
"method"
]
||
!
root
[
"from_id"
]
||
!
root
[
"pull_type"
]
)
{
values
[
"fid"
]
=
fid
;
values
[
"tid"
]
=
root
[
"tid"
].
asString
();
values
[
"candidate"
]
=
root
[
"candidate"
].
asString
();
error
=
"param error"
;
return
false
;
}
values
[
"message_id"
]
=
root
[
"message_id"
].
asString
();
values
[
"method"
]
=
root
[
"method"
].
asString
();
values
[
"from_id"
]
=
root
[
"from_id"
].
asString
();
values
[
"pull_type"
]
=
root
[
"pull_type"
].
asString
();
return
true
;
}
bool
MqttRequest
::
ParseResponse
(
std
::
string
&
response
,
std
::
string
&
path
,
ReqValue
&
values
)
bool
MqttRequest
::
ParseOfferRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
)
{
Json
::
Value
root
;
if
(
!
GetJsonValue
(
request
,
root
))
{
error
=
"json format error"
;
return
false
;
}
if
(
!
GetJsonValue
(
response
,
root
))
return
false
;
if
(
!
root
[
"path"
]
||
!
root
[
"fid"
])
return
false
;
path
=
root
[
"path"
].
asString
();
std
::
string
fid
=
root
[
"fid"
].
asString
();
if
(
path
.
compare
(
"join"
)
==
0
)
if
(
!
root
[
"message_id"
]
||
!
root
[
"method"
]
||
!
root
[
"from_id"
]
||
!
root
[
"to_id"
]
||
!
root
[
"candidate"
])
{
values
[
"ids"
]
=
root
[
"ids"
].
asString
()
;
values
[
"result"
]
=
root
[
"result"
].
asString
()
;
error
=
"param error"
;
return
false
;
}
values
[
"message_id"
]
=
root
[
"message_id"
].
asString
();
values
[
"method"
]
=
root
[
"method"
].
asString
();
values
[
"from_id"
]
=
root
[
"from_id"
].
asString
();
values
[
"to_id"
]
=
root
[
"from_id"
].
asString
();
values
[
"candidate"
]
=
root
[
"candidate"
].
asString
();
return
true
;
}
bool
MqttRequest
::
ParseAnswerRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
)
{
return
ParseOfferRequest
(
request
,
values
,
error
);
}
std
::
string
MqttRequest
::
JoinResponse
(
std
::
string
ids
,
std
::
string
result
)
bool
MqttRequest
::
ParsePublishersRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
)
{
Json
::
Value
root
;
Json
::
StreamWriterBuilder
wBuilder
;
if
(
!
GetJsonValue
(
request
,
root
))
{
error
=
"json format error"
;
return
false
;
}
root
[
"path"
]
=
"join"
;
root
[
"ids"
]
=
ids
;
root
[
"result"
]
=
result
;
if
(
!
root
[
"message_id"
]
||
!
root
[
"method"
]
||
!
root
[
"publishers"
]
||
!
root
[
"to_id"
])
{
error
=
"param error"
;
return
false
;
}
return
Json
::
writeString
(
wBuilder
,
root
);
values
[
"message_id"
]
=
root
[
"message_id"
].
asString
();
values
[
"method"
]
=
root
[
"method"
].
asString
();
values
[
"publishers_id"
]
=
root
[
"publishers_id"
].
asString
();
values
[
"to_id"
]
=
root
[
"to_id"
].
asString
();
return
true
;
}
std
::
string
MqttRequest
::
UpdateResponse
(
std
::
string
result
)
std
::
string
MqttRequest
::
BuildResponse
(
std
::
string
message_id
,
std
::
string
result_code
)
{
Json
::
Value
root
;
Json
::
StreamWriterBuilder
wBuilder
;
root
[
"
path"
]
=
"update"
;
root
[
"result
"
]
=
result
;
root
[
"
message_id"
]
=
message_id
;
root
[
"result
_code"
]
=
result_code
;
return
Json
::
writeString
(
wBuilder
,
root
);
}
...
...
thirdparty/mqtt/mqtt_request.h
View file @
40cf7643
...
...
@@ -3,6 +3,38 @@
#include <string>
#include <map>
/*
Request format
{
"message_id":"xxxx",
"method":"join",
"from_id":"xxxx",
"to_id":"xxxx"
"data":
{}
}
Response format
{
"message_id":"xxxx",
"result_code":"ok" or
"result_code":"no auth"
"data":
{}
}
client to server:
join
update
client to client:
offer/answer
server to client:
publishers
*/
namespace
offcn
{
typedef
std
::
map
<
std
::
string
,
std
::
string
>
ReqValue
;
...
...
@@ -10,14 +42,44 @@ namespace offcn
class
MqttRequest
{
public
:
static
std
::
string
JoinRequest
(
std
::
string
local_id
);
static
std
::
string
UpdateRequest
(
std
::
string
local_id
,
std
::
string
pull_type
);
static
std
::
string
Candidate
(
std
::
string
local_id
,
std
::
string
remote_id
,
std
::
string
type
,
std
::
string
candidate
);
/**
* client to server
*/
static
std
::
string
JoinRequest
(
std
::
string
message_id
,
std
::
string
from_id
);
static
std
::
string
UpdateRequest
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
pull_type
);
/**
* client to client
*/
static
std
::
string
OfferRequest
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
to_id
,
std
::
string
candidate
);
static
std
::
string
AnswerRequest
(
std
::
string
message_id
,
std
::
string
from_id
,
std
::
string
to_id
,
std
::
string
candidate
);
/**
* server to client
*/
static
std
::
string
PublishersRequest
(
std
::
string
message_id
,
std
::
string
to_id
,
std
::
string
publisher_id
);
/**
* public
*/
static
bool
ParseRequestType
(
std
::
string
request
,
std
::
string
&
type
,
std
::
string
&
error
);
static
bool
ParseRequest
(
std
::
string
&
req
,
std
::
string
&
path
,
ReqValue
&
values
);
static
bool
ParseResponse
(
std
::
string
&
response
,
std
::
string
&
path
,
ReqValue
&
values
);
/**
* client to server
*/
static
bool
ParseJoinRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
);
static
bool
ParseUpdateRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
);
/**
* client to client
*/
static
bool
ParseOfferRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
);
static
bool
ParseAnswerRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
);
/**
* server to client
*/
static
bool
ParsePublishersRequest
(
std
::
string
request
,
ReqValue
&
values
,
std
::
string
&
error
);
static
std
::
string
JoinResponse
(
std
::
string
ids
,
std
::
string
result
);
static
std
::
string
UpdateResponse
(
std
::
string
result
);
/**
* build response
*/
static
std
::
string
BuildResponse
(
std
::
string
message_id
,
std
::
string
result_code
);
};
}
\ No newline at end of file
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment