Correct RPC message processing and invocation id assignment

This commit is contained in:
Simon Fels 2016-11-26 13:05:24 +01:00
commit ac8c4a9305
3 changed files with 31 additions and 31 deletions

View file

@ -98,8 +98,9 @@ void Channel::notify_disconnected() {
pending_calls_->force_completion();
}
int Channel::next_id() {
return next_message_id_.fetch_add(1);
std::uint32_t Channel::next_id() {
static std::uint32_t next_message_id = 0;
return next_message_id++;
}
} // namespace rpc
} // namespace anbox

View file

@ -59,10 +59,9 @@ private:
std::string const& method_name,
google::protobuf::MessageLite const* request);
void send_message(const std::uint8_t &type, google::protobuf::MessageLite const& message);
int next_id();
std::uint32_t next_id();
void notify_disconnected();
std::atomic<int> next_message_id_;
std::shared_ptr<PendingCallCache> pending_calls_;
std::shared_ptr<network::MessageSender> sender_;
std::mutex write_mutex_;

View file

@ -20,6 +20,7 @@
#include "anbox/rpc/make_protobuf_object.h"
#include "anbox/rpc/constants.h"
#include "anbox/common/variable_length_array.h"
#include "anbox/utils.h"
#include "anbox_rpc.pb.h"
@ -50,39 +51,38 @@ bool MessageProcessor::process_data(const std::vector<std::uint8_t> &data) {
for (const auto &byte : data)
buffer_.push_back(byte);
buffer_.shrink_to_fit();
while (buffer_.size() > 0) {
const auto high = buffer_[0];
const auto low = buffer_[1];
size_t const message_size = (high << 8) + low;
const auto message_type = buffer_[2];
const auto high = buffer_[0];
const auto low = buffer_[1];
size_t const message_size = (high << 8) + low;
const auto message_type = buffer_[2];
// If we don't have yet all bytes for a new message return and wait
// until we have all.
if (buffer_.size() - header_size < message_size) {
WARNING("Don't have enough bytes yet (buffer %d header %d message %d)",
buffer_.size(), header_size, message_size);
break;
}
// If we don't have yet all bytes for a new message return and wait
// until we have all.
if (buffer_.size() - header_size < message_size)
return true;
if (message_type == MessageType::invocation) {
anbox::protobuf::rpc::Invocation raw_invocation;
raw_invocation.ParseFromArray(buffer_.data() + header_size, message_size);
buffer_.erase(buffer_.begin(), buffer_.begin() + header_size);
dispatch(Invocation(raw_invocation));
}
else if (message_type == MessageType::response) {
auto result = make_protobuf_object<protobuf::rpc::Result>();
result->ParseFromArray(buffer_.data() + header_size, message_size);
if (message_type == MessageType::invocation) {
anbox::protobuf::rpc::Invocation raw_invocation;
raw_invocation.ParseFromArray(buffer_.data(), message_size);
if (result->has_id())
pending_calls_->complete_response(*result);
buffer_.erase(buffer_.begin(), buffer_.begin() + message_size);
for (int n = 0; n < result->events_size(); n++)
process_event_sequence(result->events(n));
}
dispatch(Invocation(raw_invocation));
}
else if (message_type == MessageType::response) {
auto result = make_protobuf_object<protobuf::rpc::Result>();
result->ParseFromArray(buffer_.data(), message_size);
buffer_.erase(buffer_.begin(), buffer_.begin() + message_size);
if (result->has_id())
pending_calls_->complete_response(*result);
for (int n = 0; n < result->events_size(); n++)
process_event_sequence(result->events(n));
buffer_.erase(buffer_.begin(), buffer_.begin() + header_size + message_size);
}
return true;