diff --git a/src/anbox/rpc/channel.cpp b/src/anbox/rpc/channel.cpp index de52436..d613996 100644 --- a/src/anbox/rpc/channel.cpp +++ b/src/anbox/rpc/channel.cpp @@ -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 diff --git a/src/anbox/rpc/channel.h b/src/anbox/rpc/channel.h index 8b71659..3dc9e11 100644 --- a/src/anbox/rpc/channel.h +++ b/src/anbox/rpc/channel.h @@ -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 next_message_id_; std::shared_ptr pending_calls_; std::shared_ptr sender_; std::mutex write_mutex_; diff --git a/src/anbox/rpc/message_processor.cpp b/src/anbox/rpc/message_processor.cpp index 85d5d87..c3b190d 100644 --- a/src/anbox/rpc/message_processor.cpp +++ b/src/anbox/rpc/message_processor.cpp @@ -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 &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(); + 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(); - 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;