diff --git a/internal/platform/implementation/linux/ble_gatt_client.cc b/internal/platform/implementation/linux/ble_gatt_client.cc index d3fd0451..e7b02732 100644 --- a/internal/platform/implementation/linux/ble_gatt_client.cc +++ b/internal/platform/implementation/linux/ble_gatt_client.cc @@ -219,20 +219,29 @@ bool BluezGattDiscovery::InitializeKnownServices() { org::bluez::GattCharacteristic1_proxy::INTERFACE_NAME) == 1; }); - for (; chr_it != objects.cend(); chr_it++) { - const auto &[path, ifaces] = *chr_it; - const auto &properties = - ifaces.at(org::bluez::GattCharacteristic1_proxy::INTERFACE_NAME); - auto maybe_props = characteristicProperties(path, properties); - if (!maybe_props.has_value()) continue; - auto [chr_uuid, service_uuid, device_path] = *maybe_props; +for (; chr_it != objects.cend(); ++chr_it) { + const auto& [path, ifaces] = *chr_it; - discovered_characteristics_.emplace( - std::make_tuple(chr_uuid, service_uuid, device_path), path); - characteristics_properties_.emplace( - path, std::make_tuple(chr_uuid, service_uuid, device_path)); + auto iface_it = ifaces.find(org::bluez::GattCharacteristic1_proxy::INTERFACE_NAME); + if (iface_it == ifaces.end()) { + // Not a GattCharacteristic1 object (or interfaces map incomplete) -> skip + continue; } + const auto& properties = iface_it->second; + + auto maybe_props = characteristicProperties(path, properties); + if (!maybe_props.has_value()) continue; + + auto [chr_uuid, service_uuid, device_path] = *maybe_props; + + discovered_characteristics_.emplace( + std::make_tuple(chr_uuid, service_uuid, device_path), path); + + characteristics_properties_.emplace( + path, std::make_tuple(chr_uuid, service_uuid, device_path)); +} + return true; } @@ -282,6 +291,8 @@ bool BluezGattDiscovery::DiscoverServiceAndCharacteristics( }; absl::ReaderMutexLock lock(&mutex_, absl::Condition(&discovered)); + + LOG(INFO) << __func__ << ": Finished discovering gatt services and characteristics"; return !cancel.Cancelled(); } @@ -331,7 +342,7 @@ BluezGattDiscovery::GetSubscribedCharacteristic( std::optional> BluezGattDiscovery::characteristicProperties( - const sdbus::ObjectPath &path, + const sdbus::ObjectPath &char_path, const std::map &properties) { mutex_.AssertHeld(); @@ -339,31 +350,41 @@ BluezGattDiscovery::characteristicProperties( auto chr_uuid = UuidFromString(chr_uuid_str); if (!chr_uuid.has_value()) { LOG(ERROR) << ": Couldn't parse UUID '" << chr_uuid_str - << "' in characteristic " << path; + << "' in characteristic " << char_path; return std::nullopt; } const sdbus::ObjectPath &service_path = properties.at("Service"); if (cached_services_.count(service_path) == 0) { cached_services_.emplace( - path, std::make_unique(system_bus_, path)); + service_path, std::make_unique(system_bus_, service_path)); } - auto &service = cached_services_.at(service_path); - nearby::Uuid service_uuid; - try { - const std::string &service_uuid_str = service->UUID(); - auto service_uuid_maybe = UuidFromString(service_uuid_str); - if (!service_uuid_maybe.has_value()) { - LOG(ERROR) << ": Couldn't parse UUID '" << service_uuid_str - << "' in service " << service_path; - return std::nullopt; - } - service_uuid = *service_uuid_maybe; - } catch (const sdbus::Error &e) { - DBUS_LOG_PROPERTY_GET_ERROR(service, "UUID", e); +auto it = cached_services_.find(service_path); +if (it == cached_services_.end() || it->second == nullptr) { + LOG(ERROR) << ": cached_services_ missing service " << service_path + << " (from characteristic " << char_path << ")"; + return std::nullopt; +} + + LOG(INFO) << ": Found service path " << service_path + << " (from characteristic " << char_path << ")"; + auto* service = it->second.get(); // service is GattServiceClient* +nearby::Uuid service_uuid; +try { + std::string service_uuid_str = service->UUID(); // copy (safe) + auto service_uuid_maybe = UuidFromString(service_uuid_str); + + if (!service_uuid_maybe.has_value()) { + LOG(ERROR) << ": Couldn't parse UUID '" << service_uuid_str + << "' in service " << service_path; return std::nullopt; } + service_uuid = *service_uuid_maybe; +} catch (const sdbus::Error &e) { + DBUS_LOG_PROPERTY_GET_ERROR(service, "UUID", e); + return std::nullopt; +} sdbus::ObjectPath device_path; try { @@ -417,7 +438,14 @@ void BluezGattDiscovery::onInterfacesRemoved( if (chr_it != end) { absl::MutexLock lock(&mutex_); { - auto &props = characteristics_properties_.at(objectPath); + auto it = characteristics_properties_.find(objectPath); + if (it == characteristics_properties_.end()) { + // Not tracked / already removed / never added. + // return; // or just `break;` / `continue;` depending on your context + return; + } + + auto &props = it->second; discovered_characteristics_.erase(props); } characteristics_properties_.erase(objectPath); diff --git a/internal/platform/implementation/linux/ble_gatt_client.h b/internal/platform/implementation/linux/ble_gatt_client.h index f86f9970..452b79a0 100644 --- a/internal/platform/implementation/linux/ble_gatt_client.h +++ b/internal/platform/implementation/linux/ble_gatt_client.h @@ -76,7 +76,7 @@ class BluezGattDiscovery final : public bluez::BluezObjectManager { private: std::optional> characteristicProperties( - const sdbus::ObjectPath &path, + const sdbus::ObjectPath &char_path, const std::map &properties) ABSL_EXCLUSIVE_LOCKS_REQUIRED(mutex_); diff --git a/internal/platform/implementation/linux/ble_gatt_server.cc b/internal/platform/implementation/linux/ble_gatt_server.cc index 0c2089ea..1c94361c 100644 --- a/internal/platform/implementation/linux/ble_gatt_server.cc +++ b/internal/platform/implementation/linux/ble_gatt_server.cc @@ -61,7 +61,9 @@ GattServer::CreateCharacteristic( if (service->AddCharacteristic(service_uuid, characteristic_uuid, permission, property)) { try { - LOG(INFO)<< __func__ << ": Registering service on gattmanager"; + LOG(INFO)<< __func__ << ": Registering service on gattmanager with characteristic_uuid: " + << std::string(characteristic_uuid) << " and service_uuid: " << std::string(service_uuid); + gatt_manager_ -> RegisterApplication(gatt_service_root_object_manager -> getObjectPath(), {}); } catch (const sdbus::Error& e) { LOG(ERROR) @@ -75,6 +77,7 @@ GattServer::CreateCharacteristic( services_.insert({service_uuid, std::move(service)}); + api::ble_v2::GattCharacteristic characteristic{ characteristic_uuid, service_uuid, permission, property}; return characteristic; diff --git a/internal/platform/implementation/linux/ble_v2_medium.cc b/internal/platform/implementation/linux/ble_v2_medium.cc index 518f44c7..2be204b0 100644 --- a/internal/platform/implementation/linux/ble_v2_medium.cc +++ b/internal/platform/implementation/linux/ble_v2_medium.cc @@ -18,6 +18,8 @@ #include #include +#include +#include #include #include #include @@ -29,11 +31,16 @@ // #include "internal/platform/implementation/linux/ble_gatt_server.h" #include "internal/platform/implementation/linux/ble_v2_medium.h" +#include "ble_gatt_client.h" #include "ble_gatt_server.h" +#include "ble_l2cap_server_socket.h" +#include "ble_l2cap_socket.h" +#include "bluez_le_bearer_client.h" #include "internal/platform/implementation/linux/bluetooth_classic_device.h" #include "internal/platform/implementation/linux/bluetooth_devices.h" #include "internal/platform/implementation/linux/bluez.h" #include "internal/platform/mac_address.h" +#include "internal/platform/prng.h" #include "absl/types/span.h" #include "internal/platform/implementation/linux/bluez_advertisement_monitor.h" #include "internal/platform/implementation/linux/bluez_advertisement_monitor_manager.h" @@ -219,6 +226,7 @@ BleV2Medium::StartAdvertising( std::unique_ptr BleV2Medium::StartGattServer( api::ble_v2::ServerGattConnectionCallback callback) { (void)callback; + return nullptr; return std::make_unique( *system_bus_, adapter_, devices_,std::move(callback) @@ -236,13 +244,12 @@ std::unique_ptr BleV2Medium::ConnectToGattServer( << ": GATT client connection is not supported on Linux yet."; return nullptr; } - + // This is supposed to be for a socket on top of Weave protocol. std::unique_ptr BleV2Medium::Connect( const std::string &service_id, api::ble_v2::TxPowerLevel tx_power_level, api::ble_v2::BlePeripheral::UniqueId peripheral_id, CancellationFlag *cancellation_flag) { - auto device = devices_ -> get_device_by_unique_id(peripheral_id); - LOG(INFO) << __func__ << ": Resolved device with address " << device -> GetMacAddress(); + LOG(INFO) << __func__ << ": Not implemented on linux "; return nullptr; } @@ -498,8 +505,16 @@ std::unique_ptr BleV2Medium::OpenServerSocket( std::unique_ptr BleV2Medium::OpenL2capServerSocket(const std::string &service_id) { - LOG(WARNING) << __func__ << ": L2CAP server sockets not implemented on Linux"; - return nullptr; + LOG(INFO) << __func__ << ": Opening L2CAP server socket for service " + << service_id; + + Prng prng; + auto psm = 0x80 + (prng.NextUint32() % 0x80); + auto server_socket = std::make_unique(psm); + + LOG(INFO) << __func__ << ": L2CAP server socket created with PSM: " + << server_socket->GetPSM(); + return server_socket; } // std::unique_ptr BleV2Medium::Connect( @@ -515,8 +530,47 @@ std::unique_ptr BleV2Medium::ConnectOverL2cap( api::ble_v2::TxPowerLevel tx_power_level, api::ble_v2::BlePeripheral::UniqueId peripheral_id, CancellationFlag *cancellation_flag) { - LOG(WARNING) << __func__ << ": L2CAP socket connections not implemented on Linux"; - return nullptr; + auto device = devices_->get_device_by_unique_id(peripheral_id); + if (!device) { + LOG(ERROR) << __func__ << ": Failed to find device with unique ID " + << peripheral_id; + return nullptr; + } + + LOG(INFO) << __func__ << ": Connecting to L2CAP PSM " << psm + << " on device " << device->GetMacAddress(); + + + int fd = socket(AF_BLUETOOTH, SOCK_SEQPACKET, BTPROTO_L2CAP); + if (fd < 0) { + LOG(ERROR) << __func__ << ": Failed to create L2CAP socket: " + << std::strerror(errno); + return nullptr; + } + + struct sockaddr_l2 addr; + std::memset(&addr, 0, sizeof(addr)); + addr.l2_family = AF_BLUETOOTH; + addr.l2_psm = htobs(psm); + addr.l2_cid = 0; + addr.l2_bdaddr_type = BDADDR_LE_PUBLIC; + + std::string mac_addr = device->GetMacAddress(); + if (str2ba(mac_addr.c_str(), &addr.l2_bdaddr) < 0) { + LOG(ERROR) << __func__ << ": Invalid Bluetooth address: " << mac_addr; + close(fd); + return nullptr; + } + + if (connect(fd, (struct sockaddr*)&addr, sizeof(addr)) < 0) { + LOG(ERROR) << __func__ << ": Failed to connect to L2CAP socket: " + << std::strerror(errno); + close(fd); + return nullptr; + } + + LOG(INFO) << __func__ << ": Successfully connected to L2CAP socket"; + return std::make_unique(fd, peripheral_id); } bool BleV2Medium::StartMultipleServicesScanning( diff --git a/internal/platform/implementation/linux/ble_v2_medium.h b/internal/platform/implementation/linux/ble_v2_medium.h index a1187a3f..fdbcce1b 100644 --- a/internal/platform/implementation/linux/ble_v2_medium.h +++ b/internal/platform/implementation/linux/ble_v2_medium.h @@ -26,7 +26,10 @@ #include "absl/synchronization/mutex.h" #include "internal/platform/implementation/ble_v2.h" // #include "internal/platform/implementation/linux/ble_gatt_client.h" +#include "ble_gatt_client.h" #include "bluez_gatt_manager.h" +#include "internal/platform/implementation/linux/ble_l2cap_server_socket.h" +#include "internal/platform/implementation/linux/ble_l2cap_socket.h" #include "internal/platform/implementation/linux/ble_v2_server_socket.h" #include "internal/platform/implementation/linux/bluetooth_adapter.h" #include "internal/platform/implementation/linux/bluetooth_devices.h" @@ -141,7 +144,7 @@ class BleV2Medium final : public api::ble_v2::BleMedium { BluetoothAdapter adapter_; ObserverList observers_ = {}; std::shared_ptr devices_; - // std::shared_ptr gatt_discovery_; + std::shared_ptr gatt_discovery_; std::unique_ptr root_object_manager_; std::unique_ptr adv_monitor_manager_; diff --git a/internal/platform/implementation/linux/ble_v2_server_socket.cc b/internal/platform/implementation/linux/ble_v2_server_socket.cc new file mode 100644 index 00000000..240de1e0 --- /dev/null +++ b/internal/platform/implementation/linux/ble_v2_server_socket.cc @@ -0,0 +1,86 @@ +// Copyright 2024 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#include "internal/platform/implementation/linux/ble_v2_server_socket.h" + +#include + +#include "absl/synchronization/mutex.h" +#include "internal/platform/exception.h" +#include "internal/platform/implementation/ble_v2.h" +#include "internal/platform/implementation/linux/ble_v2_socket.h" +#include "internal/platform/logging.h" + +namespace nearby { +namespace linux { + +std::unique_ptr BleV2ServerSocket::Accept() { + absl::MutexLock lock(&mutex_); + LOG(INFO) << "BleV2ServerSocket::Accept waiting for connection"; + + while (!closed_ && pending_sockets_.empty()) { + cond_.Wait(&mutex_); + } + + if (closed_) { + LOG(INFO) << "BleV2ServerSocket::Accept socket is closed"; + return nullptr; + } + + std::unique_ptr socket = std::move(pending_sockets_.front()); + pending_sockets_.pop_front(); + + LOG(INFO) << "BleV2ServerSocket::Accept accepted connection"; + return socket; +} + +Exception BleV2ServerSocket::Close() { + absl::MutexLock lock(&mutex_); + LOG(INFO) << "BleV2ServerSocket::Close for service " << service_id_; + + if (closed_) { + return {Exception::kSuccess}; + } + + closed_ = true; + + // Close all pending sockets + for (auto& socket : pending_sockets_) { + if (socket) { + socket->Close(); + } + } + pending_sockets_.clear(); + + cond_.SignalAll(); + + return {Exception::kSuccess}; +} + +void BleV2ServerSocket::AddPendingSocket(std::unique_ptr socket) { + absl::MutexLock lock(&mutex_); + if (closed_) { + LOG(WARNING) + << "BleV2ServerSocket::AddPendingSocket socket is closed"; + return; + } + + pending_sockets_.push_back(std::move(socket)); + cond_.SignalAll(); + LOG(INFO) << "BleV2ServerSocket::AddPendingSocket added socket, " + << "pending count: " << pending_sockets_.size(); +} + +} // namespace linux +} // namespace nearby diff --git a/internal/platform/implementation/linux/ble_v2_socket.cc b/internal/platform/implementation/linux/ble_v2_socket.cc new file mode 100644 index 00000000..b58303e2 --- /dev/null +++ b/internal/platform/implementation/linux/ble_v2_socket.cc @@ -0,0 +1,264 @@ +// Copyright 2024 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#include "internal/platform/implementation/linux/ble_v2_socket.h" + +#include + +#include "absl/synchronization/mutex.h" +#include "internal/platform/byte_array.h" +#include "internal/platform/exception.h" +#include "internal/platform/logging.h" + +namespace nearby { +namespace linux { + +InputStream& BleV2Socket::GetInputStream() { return input_stream_; } + +OutputStream& BleV2Socket::GetOutputStream() { return output_stream_; } + +Exception BleV2Socket::Close() { + absl::MutexLock lock(&mutex_); + if (closed_) { + return {Exception::kSuccess}; + } + closed_ = true; + + // Close streams + input_stream_.NotifyClose(); + output_stream_.Close(); + + // Cleanup GATT resources + if (gatt_client_) { + LOG(INFO) << "Disconnecting GATT client for peripheral " + << peripheral_id_; + gatt_client_->Disconnect(); + gatt_client_.reset(); + } + + if (gatt_server_) { + LOG(INFO) << "Stopping GATT server for peripheral " + << peripheral_id_; + // Server cleanup is handled by the server itself + gatt_server_.reset(); + } + + return {Exception::kSuccess}; +} + +bool BleV2Socket::IsClosed() const { + absl::MutexLock lock(&mutex_); + return closed_; +} + +// BleInputStream implementation +ExceptionOr BleV2Socket::BleInputStream::Read(std::int64_t size) { + absl::MutexLock lock(&mutex_); + + while (buffer_.Empty() && !closed_) { + cond_.Wait(&mutex_); + } + + if (closed_ && buffer_.Empty()) { + return ExceptionOr(Exception::kIo); + } + + if (size < 0 || static_cast(size) >= buffer_.size()) { + ByteArray result = buffer_; + buffer_ = ByteArray(); + return ExceptionOr(result); + } + + ByteArray result(buffer_.data(), size); + buffer_ = ByteArray(buffer_.data() + size, buffer_.size() - size); + return ExceptionOr(result); +} + +Exception BleV2Socket::BleInputStream::Close() { + absl::MutexLock lock(&mutex_); + if (closed_) { + return {Exception::kSuccess}; + } + closed_ = true; + cond_.SignalAll(); + return {Exception::kSuccess}; +} + +void BleV2Socket::BleInputStream::ReceiveData(const ByteArray& data) { + absl::MutexLock lock(&mutex_); + if (closed_) { + return; + } + + if (buffer_.Empty()) { + buffer_ = data; + } else { + ByteArray combined(buffer_.size() + data.size()); + std::memcpy(combined.data(), buffer_.data(), buffer_.size()); + std::memcpy(combined.data() + buffer_.size(), data.data(), data.size()); + buffer_ = std::move(combined); + } + cond_.SignalAll(); +} + +void BleV2Socket::BleInputStream::NotifyClose() { + absl::MutexLock lock(&mutex_); + closed_ = true; + cond_.SignalAll(); +} + +// BleOutputStream implementation +Exception BleV2Socket::BleOutputStream::Write(const ByteArray& data) { + absl::MutexLock lock(&mutex_); + if (closed_) { + return {Exception::kIo}; + } + + if (!write_callback_) { + LOG(WARNING) << "BleOutputStream: No write callback set"; + return {Exception::kIo}; + } + + bool success = write_callback_(data); + return {success ? Exception::kSuccess : Exception::kIo}; +} + +Exception BleV2Socket::BleOutputStream::Flush() { + return {Exception::kSuccess}; +} + +Exception BleV2Socket::BleOutputStream::Close() { + absl::MutexLock lock(&mutex_); + if (closed_) { + return {Exception::kSuccess}; + } + closed_ = true; + write_callback_ = nullptr; + return {Exception::kSuccess}; +} + +void BleV2Socket::BleOutputStream::SetWriteCallback(WriteCallback callback) { + absl::MutexLock lock(&mutex_); + write_callback_ = std::move(callback); +} + +void BleV2Socket::SetGattServer( + std::unique_ptr gatt_server, + const api::ble_v2::GattCharacteristic& rx_char, + const api::ble_v2::GattCharacteristic& tx_char) { + absl::MutexLock lock(&mutex_); + + if (closed_) { + LOG(WARNING) << "Cannot set GATT server on closed socket"; + return; + } + + gatt_server_ = std::move(gatt_server); + rx_char_ = rx_char; + tx_char_ = tx_char; + + // Set up write callback to use GATT server notifications + output_stream_.SetWriteCallback( + [this](const ByteArray& data) -> bool { + absl::MutexLock lock(&mutex_); + if (!gatt_server_) { + LOG(ERROR) << "GATT server not available for write"; + return false; + } + + if (closed_) { + LOG(WARNING) << "Socket is closed, cannot write"; + return false; + } + + // Notify remote device via TX characteristic + absl::Status status = gatt_server_->NotifyCharacteristicChanged( + tx_char_, /*confirm=*/false, data); + + if (!status.ok()) { + LOG(WARNING) << "Failed to notify TX characteristic: " + << status.message(); + return false; + } + return true; + }); + + LOG(INFO) << "BLE socket configured with GATT server, RX: " + << std::string(rx_char.uuid) + << ", TX: " << std::string(tx_char.uuid); +} + +void BleV2Socket::SetGattClient( + std::unique_ptr gatt_client, + const api::ble_v2::GattCharacteristic& rx_char, + const api::ble_v2::GattCharacteristic& tx_char) { + absl::MutexLock lock(&mutex_); + + if (closed_) { + LOG(WARNING) << "Cannot set GATT client on closed socket"; + return; + } + + gatt_client_ = std::move(gatt_client); + rx_char_ = rx_char; + tx_char_ = tx_char; + + // Set up write callback to use GATT client writes + output_stream_.SetWriteCallback( + [this](const ByteArray& data) -> bool { + absl::MutexLock lock(&mutex_); + if (!gatt_client_) { + LOG(ERROR) << "GATT client not available for write"; + return false; + } + + if (closed_) { + LOG(WARNING) << "Socket is closed, cannot write"; + return false; + } + + // Write to TX characteristic on remote device + std::string data_str(data.data(), data.size()); + bool success = gatt_client_->WriteCharacteristic( + tx_char_, data_str, + api::ble_v2::GattClient::WriteType::kWithoutResponse); + + if (!success) { + LOG(WARNING) << "Failed to write to TX characteristic"; + return false; + } + return true; + }); + + // Subscribe to RX characteristic to receive data + bool subscribed = gatt_client_->SetCharacteristicSubscription( + rx_char_, /*enable=*/true, + [this](absl::string_view value) { + if (!IsClosed()) { + ByteArray data(value.data(), value.size()); + input_stream_.ReceiveData(data); + } + }); + + if (!subscribed) { + LOG(ERROR) << "Failed to subscribe to RX characteristic"; + } + + LOG(INFO) << "BLE socket configured with GATT client, RX: " + << std::string(rx_char.uuid) + << ", TX: " << std::string(tx_char.uuid); +} + +} // namespace linux +} // namespace nearby diff --git a/internal/platform/implementation/linux/ble_v2_socket.h b/internal/platform/implementation/linux/ble_v2_socket.h new file mode 100644 index 00000000..82ea93c2 --- /dev/null +++ b/internal/platform/implementation/linux/ble_v2_socket.h @@ -0,0 +1,138 @@ +// Copyright 2024 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#ifndef PLATFORM_IMPL_LINUX_API_BLE_V2_SOCKET_H_ +#define PLATFORM_IMPL_LINUX_API_BLE_V2_SOCKET_H_ + +#include +#include + +#include "absl/synchronization/mutex.h" +#include "internal/platform/byte_array.h" +#include "internal/platform/exception.h" +#include "internal/platform/implementation/ble_v2.h" +#include "internal/platform/input_stream.h" +#include "internal/platform/output_stream.h" + +namespace nearby { +namespace linux { + +// BLE v2 Socket implementation using GATT characteristics for data transfer +// +// Data flow: +// - Server side: +// * RX characteristic: Remote writes -> our InputStream reads +// * TX characteristic: Our OutputStream writes -> remote reads via notifications +// - Client side: +// * TX characteristic: Our OutputStream writes -> remote reads +// * RX characteristic: Remote writes (notifications) -> our InputStream reads +class BleV2Socket : public api::ble_v2::BleSocket { + public: + BleV2Socket() = default; + explicit BleV2Socket(api::ble_v2::BlePeripheral::UniqueId peripheral_id) + : peripheral_id_(peripheral_id) {} + ~BleV2Socket() override { Close(); } + + InputStream& GetInputStream() override; + OutputStream& GetOutputStream() override; + Exception Close() override ABSL_LOCKS_EXCLUDED(mutex_); + + api::ble_v2::BlePeripheral::UniqueId GetRemotePeripheralId() override { + return peripheral_id_; + } + + bool IsClosed() const ABSL_LOCKS_EXCLUDED(mutex_); + + // GATT integration: Allow external code to feed data to input stream + // Called when remote device writes to RX characteristic + void ReceiveData(const ByteArray& data) { input_stream_.ReceiveData(data); } + + // GATT integration: Set callback for output stream writes + // Callback should write to TX characteristic (notify remote) + void SetWriteCallback( + absl::AnyInvocable callback) { + output_stream_.SetWriteCallback(std::move(callback)); + } + + // Set the GATT server and characteristics for server-side socket + void SetGattServer(std::unique_ptr gatt_server, + const api::ble_v2::GattCharacteristic& rx_char, + const api::ble_v2::GattCharacteristic& tx_char) + ABSL_LOCKS_EXCLUDED(mutex_); + + // Set the GATT client and characteristics for client-side socket + void SetGattClient(std::unique_ptr gatt_client, + const api::ble_v2::GattCharacteristic& rx_char, + const api::ble_v2::GattCharacteristic& tx_char) + ABSL_LOCKS_EXCLUDED(mutex_); + + private: + class BleInputStream : public InputStream { + public: + BleInputStream() = default; + ~BleInputStream() override = default; + + ExceptionOr Read(std::int64_t size) override + ABSL_LOCKS_EXCLUDED(mutex_); + Exception Close() override ABSL_LOCKS_EXCLUDED(mutex_); + + void ReceiveData(const ByteArray& data) ABSL_LOCKS_EXCLUDED(mutex_); + void NotifyClose() ABSL_LOCKS_EXCLUDED(mutex_); + + private: + absl::Mutex mutex_; + absl::CondVar cond_; + ByteArray buffer_ ABSL_GUARDED_BY(mutex_); + bool closed_ ABSL_GUARDED_BY(mutex_) = false; + }; + + class BleOutputStream : public OutputStream { + public: + BleOutputStream() = default; + ~BleOutputStream() override = default; + + Exception Write(const ByteArray& data) override + ABSL_LOCKS_EXCLUDED(mutex_); + Exception Flush() override; + Exception Close() override ABSL_LOCKS_EXCLUDED(mutex_); + + using WriteCallback = absl::AnyInvocable; + void SetWriteCallback(WriteCallback callback) + ABSL_LOCKS_EXCLUDED(mutex_); + + private: + absl::Mutex mutex_; + WriteCallback write_callback_ ABSL_GUARDED_BY(mutex_); + bool closed_ ABSL_GUARDED_BY(mutex_) = false; + }; + + mutable absl::Mutex mutex_; + bool closed_ ABSL_GUARDED_BY(mutex_) = false; + BleInputStream input_stream_; + BleOutputStream output_stream_; + api::ble_v2::BlePeripheral::UniqueId peripheral_id_ = 0; + + // GATT resources (only one of these will be set) + std::unique_ptr gatt_server_ ABSL_GUARDED_BY(mutex_); + std::unique_ptr gatt_client_ ABSL_GUARDED_BY(mutex_); + + // Characteristics for data transfer + api::ble_v2::GattCharacteristic rx_char_ ABSL_GUARDED_BY(mutex_); + api::ble_v2::GattCharacteristic tx_char_ ABSL_GUARDED_BY(mutex_); +}; + +} // namespace linux +} // namespace nearby + +#endif // PLATFORM_IMPL_LINUX_API_BLE_V2_SOCKET_H_ diff --git a/internal/platform/implementation/linux/bluetooth_classic_device.cc b/internal/platform/implementation/linux/bluetooth_classic_device.cc index 9805994c..007b1909 100644 --- a/internal/platform/implementation/linux/bluetooth_classic_device.cc +++ b/internal/platform/implementation/linux/bluetooth_classic_device.cc @@ -107,7 +107,18 @@ bool BluetoothDevice::ConnectToProfile(absl::string_view service_uuid) { } } -MonitoredBluetoothDevice::MonitoredBluetoothDevice( + bool BluetoothDevice::Connect() { + auto device = device_; + if (device == nullptr) return false; + try { + device->Connect(); + return true; + } catch (const sdbus::Error &e) { + DBUS_LOG_METHOD_CALL_ERROR(device, "Connect", e); + return false; + } +} + MonitoredBluetoothDevice::MonitoredBluetoothDevice( std::shared_ptr system_bus, std::shared_ptr device, ObserverList &observers) @@ -152,7 +163,7 @@ void MonitoredBluetoothDevice::onPropertiesChanged( observer->DeviceConnectedStateChanged(*this, it->second); } } else if ( it -> first == "ServicesResolved"){ - LOG(INFO) << ": ServicesResolved"; + LOG(INFO) << ": ServicesResolved :" << std::string(it->second); }else if (it->first == bluez::DEVICE_NAME) { auto callback = GetDiscoveryCallback(); if (callback != nullptr && callback->device_name_changed_cb != nullptr) diff --git a/internal/platform/implementation/linux/bluetooth_classic_device.h b/internal/platform/implementation/linux/bluetooth_classic_device.h index d4310644..27d0b7e7 100644 --- a/internal/platform/implementation/linux/bluetooth_classic_device.h +++ b/internal/platform/implementation/linux/bluetooth_classic_device.h @@ -111,9 +111,11 @@ class BluetoothDevice : public api::BluetoothDevice { } bool ConnectToProfile(absl::string_view service_uuid); + bool Connect(); void MarkLost() { lost_ = true; } void UnmarkLost() { lost_ = false; } bool Lost() const { return lost_; } + sdbus::ObjectPath GetObjectPath() {return device_->getObjectPath();} private: UniqueId unique_id_; diff --git a/internal/platform/implementation/linux/bluez_le_bearer_client.h b/internal/platform/implementation/linux/bluez_le_bearer_client.h new file mode 100644 index 00000000..571bcf6d --- /dev/null +++ b/internal/platform/implementation/linux/bluez_le_bearer_client.h @@ -0,0 +1,79 @@ +// Copyright 2024 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#ifndef PLATFORM_IMPL_LINUX_LE_BEARER_CLIENT_H_ +#define PLATFORM_IMPL_LINUX_LE_BEARER_CLIENT_H_ + +#include +#include + +#include +#include +#include + +#include "absl/functional/any_invocable.h" +#include "absl/synchronization/mutex.h" +#include "internal/platform/implementation/linux/generated/dbus/bluez/le_bearer_client.h" + +namespace nearby { +namespace linux { +namespace bluez { + +class LEBearerClient + : public sdbus::ProxyInterfaces { + public: + LEBearerClient(std::shared_ptr system_bus, + sdbus::ObjectPath bearer_path) + : ProxyInterfaces(*system_bus, "org.bluez", std::move(bearer_path)), + system_bus_(std::move(system_bus)) { + registerProxy(); + } + ~LEBearerClient() { unregisterProxy(); } + + void SetDisconnectedCallback( + absl::AnyInvocable cb) + ABSL_LOCKS_EXCLUDED(disconnect_callback_lock_) { + absl::MutexLock l(&disconnect_callback_lock_); + on_disconnected_cb_ = std::move(cb); + } + + void ResetDisconnectedCallback() + ABSL_LOCKS_EXCLUDED(disconnect_callback_lock_) { + absl::MutexLock l(&disconnect_callback_lock_); + on_disconnected_cb_ = nullptr; + } + + protected: + void onDisconnected(const std::string& reason, + const std::string& message) override + ABSL_LOCKS_EXCLUDED(disconnect_callback_lock_) { + absl::ReaderMutexLock l(&disconnect_callback_lock_); + if (on_disconnected_cb_ != nullptr) { + on_disconnected_cb_(reason, message); + } + } + + private: + std::shared_ptr system_bus_; + absl::Mutex disconnect_callback_lock_; + absl::AnyInvocable + on_disconnected_cb_ ABSL_GUARDED_BY(disconnect_callback_lock_) = nullptr; +}; + +} // namespace bluez +} // namespace linux +} // namespace nearby + +#endif // PLATFORM_IMPL_LINUX_LE_BEARER_CLIENT_H_