123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250 |
- // Copyright 2021 The Chromium Authors. All rights reserved.
- // Use of this source code is governed by a BSD-style license that can be
- // found in the LICENSE file.
- #include "ash/quick_pair/message_stream/message_stream.h"
- #include "ash/quick_pair/common/fast_pair/fast_pair_metrics.h"
- #include "ash/quick_pair/common/logging.h"
- #include "ash/services/quick_pair/quick_pair_process.h"
- #include "ash/services/quick_pair/quick_pair_process_manager.h"
- #include "base/bind.h"
- #include "base/strings/string_number_conversions.h"
- #include "device/bluetooth/bluetooth_socket.h"
- #include "net/base/io_buffer.h"
- namespace {
- constexpr int kMaxBufferSize = 4096;
- constexpr int kMaxRetryCount = 10;
- constexpr int kMessageStorageCapacity = 1000;
- } // namespace
- namespace ash {
- namespace quick_pair {
- MessageStream::MessageStream(const std::string& device_address,
- scoped_refptr<device::BluetoothSocket> socket)
- : device_address_(device_address), socket_(socket) {
- Receive();
- }
- MessageStream::~MessageStream() {
- if (socket_.get())
- socket_->Disconnect(base::DoNothing());
- // Notify observers for lifetime management
- for (auto& obs : observers_)
- obs.OnMessageStreamDestroyed(device_address_);
- }
- void MessageStream::AddObserver(Observer* observer) {
- observers_.AddObserver(observer);
- }
- void MessageStream::RemoveObserver(Observer* observer) {
- observers_.RemoveObserver(observer);
- }
- void MessageStream::Receive() {
- if (receive_retry_counter_ == kMaxRetryCount) {
- QP_LOG(WARNING)
- << __func__
- << ": Failed to receive or parse data from socket more than "
- << kMaxRetryCount << " times.";
- if (socket_.get()) {
- socket_->Disconnect(base::BindOnce(&MessageStream::OnSocketDisconnected,
- weak_ptr_factory_.GetWeakPtr()));
- }
- return;
- }
- // Retry receiving data.
- receive_retry_counter_++;
- socket_->Receive(/*buffer_size=*/kMaxBufferSize,
- base::BindOnce(&MessageStream::ReceiveDataSuccess,
- weak_ptr_factory_.GetWeakPtr()),
- base::BindOnce(&MessageStream::ReceiveDataError,
- weak_ptr_factory_.GetWeakPtr()));
- }
- void MessageStream::ReceiveDataSuccess(int buffer_size,
- scoped_refptr<net::IOBuffer> io_buffer) {
- RecordMessageStreamReceiveResult(/*success=*/true);
- receive_retry_counter_ = 0;
- if (!io_buffer->data()) {
- Receive();
- return;
- }
- std::vector<uint8_t> message_bytes(buffer_size);
- for (int i = 0; i < buffer_size; i++) {
- char* c = io_buffer->data() + i;
- message_bytes[i] = static_cast<uint8_t>(*c);
- }
- quick_pair_process::ParseMessageStreamMessages(
- std::move(message_bytes),
- base::BindOnce(&MessageStream::ParseMessageStreamSuccess,
- weak_ptr_factory_.GetWeakPtr()),
- base::BindOnce(&MessageStream::OnUtilityProcessStopped,
- weak_ptr_factory_.GetWeakPtr()));
- }
- void MessageStream::ReceiveDataError(device::BluetoothSocket::ErrorReason error,
- const std::string& error_message) {
- QP_LOG(INFO) << __func__ << ": Error: " << error_message;
- RecordMessageStreamReceiveResult(/*success=*/false);
- RecordMessageStreamReceiveError(error);
- if (error == device::BluetoothSocket::ErrorReason::kDisconnected) {
- OnSocketDisconnected();
- return;
- }
- Receive();
- }
- void MessageStream::Disconnect(base::OnceClosure on_disconnect_callback) {
- // If we already have disconnected the socket, then we can run the callback.
- // This can happen since the socket might have disconnected previously but
- // we kept the MessageStream instance alive to preserve messages from the
- // corresponding device.
- if (!socket_.get()) {
- std::move(on_disconnect_callback).Run();
- return;
- }
- socket_->Disconnect(base::BindOnce(
- &MessageStream::OnSocketDisconnectedWithCallback,
- weak_ptr_factory_.GetWeakPtr(), std::move(on_disconnect_callback)));
- }
- void MessageStream::OnSocketDisconnected() {
- for (auto& obs : observers_)
- obs.OnDisconnected(device_address_);
- }
- void MessageStream::OnSocketDisconnectedWithCallback(
- base::OnceClosure on_disconnect_callback) {
- OnSocketDisconnected();
- std::move(on_disconnect_callback).Run();
- }
- void MessageStream::ParseMessageStreamSuccess(
- std::vector<mojom::MessageStreamMessagePtr> messages) {
- QP_LOG(VERBOSE) << __func__;
- if (messages.empty()) {
- Receive();
- return;
- }
- // Store messages and notify observers.
- for (size_t i = 0; i < messages.size(); ++i) {
- if (messages_.size() == kMessageStorageCapacity)
- messages_.pop_front();
- messages_.push_back(std::move(messages[i]));
- NotifyObservers(messages_.back());
- }
- // Attempt to receive new messages from socket.
- Receive();
- }
- void MessageStream::NotifyObservers(
- const mojom::MessageStreamMessagePtr& message) {
- if (message->is_model_id()) {
- for (auto& obs : observers_)
- obs.OnModelIdMessage(device_address_, message->get_model_id());
- return;
- }
- if (message->is_ble_address_update()) {
- for (auto& obs : observers_)
- obs.OnBleAddressUpdateMessage(device_address_,
- message->get_ble_address_update());
- return;
- }
- if (message->is_battery_update()) {
- for (auto& obs : observers_)
- obs.OnBatteryUpdateMessage(device_address_,
- std::move(message->get_battery_update()));
- return;
- }
- if (message->is_remaining_battery_time()) {
- for (auto& obs : observers_)
- obs.OnRemainingBatteryTimeMessage(device_address_,
- message->get_remaining_battery_time());
- return;
- }
- if (message->is_enable_silence_mode()) {
- for (auto& obs : observers_)
- obs.OnEnableSilenceModeMessage(device_address_,
- message->get_enable_silence_mode());
- return;
- }
- if (message->is_companion_app_log_buffer_full()) {
- for (auto& obs : observers_)
- obs.OnCompanionAppLogBufferFullMessage(device_address_);
- return;
- }
- if (message->is_active_components_byte()) {
- for (auto& obs : observers_)
- obs.OnActiveComponentsMessage(device_address_,
- message->get_active_components_byte());
- return;
- }
- if (message->is_ring_device_event()) {
- for (auto& obs : observers_)
- obs.OnRingDeviceMessage(device_address_,
- std::move(message->get_ring_device_event()));
- return;
- }
- if (message->is_acknowledgement()) {
- for (auto& obs : observers_)
- obs.OnAcknowledgementMessage(device_address_,
- std::move(message->get_acknowledgement()));
- return;
- }
- if (message->is_sdk_version()) {
- for (auto& obs : observers_)
- obs.OnAndroidSdkVersionMessage(device_address_,
- message->get_sdk_version());
- return;
- }
- }
- void MessageStream::OnUtilityProcessStopped(
- QuickPairProcessManager::ShutdownReason shutdown_reason) {
- QP_LOG(INFO) << __func__ << ": Error: " << shutdown_reason;
- receive_retry_counter_++;
- Receive();
- }
- } // namespace quick_pair
- } // namespace ash
|