/* * blebus.cpp — BLEBus v1.1.0 * Pure C++17 BLE L2CAP IPC library. * Copyright (c) 2026 Are Bjørby * SPDX-License-Identifier: MIT */ #include "blebus.hpp" #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include namespace blebus { namespace { /* ── Wire frame ──────────────────────────────────────────────────────── * * [4 bytes: uint32_t payload length, network byte order][payload] * * Same framing as sockbus for wire compatibility. */ constexpr size_t FRAME_HDR_SIZE = 4; /* ── Limits ──────────────────────────────────────────────────────────── */ constexpr int MAX_BROKERS = 4; /* Custom 128-bit service UUID for BLE Antheos transport discovery. */ constexpr const char* BLEBUS_SERVICE_UUID = "6a8f1d3e-7b2c-4e0a-9f5d-8c3a2b1e0d4f"; /* ── Unix socket path ────────────────────────────────────────────────── */ std::string sock_path(std::string_view name) { return "/tmp/blebus_" + std::string(name) + ".sock"; } /* ── Frame send (SEQPACKET — single atomic send) ─────────────────── */ ssize_t send_frame(int fd, const uint8_t* data, size_t len) { size_t total = FRAME_HDR_SIZE + len; uint8_t sbuf[4096]; std::unique_ptr heap; uint8_t* buf = sbuf; if (total > sizeof(sbuf)) { heap = std::make_unique(total); buf = heap.get(); } uint32_t hdr = htonl(static_cast(len)); std::memcpy(buf, &hdr, FRAME_HDR_SIZE); if (len > 0) std::memcpy(buf + FRAME_HDR_SIZE, data, len); ssize_t ret = ::send(fd, buf, total, MSG_NOSIGNAL); return (ret < 0) ? -1 : static_cast(len); } /* ── Extract payload from received frame ─────────────────────────── */ size_t extract_payload(const uint8_t* frame, size_t frame_len, const uint8_t** payload) { if (frame_len < FRAME_HDR_SIZE) return 0; uint32_t plen; std::memcpy(&plen, frame, sizeof(plen)); plen = ntohl(plen); if (frame_len < FRAME_HDR_SIZE + plen) return 0; *payload = frame + FRAME_HDR_SIZE; return plen; } /* ═══════════════════════════════════════════════════════════════════════ * BLE Advertising via BlueZ D-Bus (sd-bus) * ═══════════════════════════════════════════════════════════════════════ */ struct Advertiser { sd_bus* bus = nullptr; sd_bus_slot* slot = nullptr; std::string local_name; bool registered = false; int fd() const { return bus ? sd_bus_get_fd(bus) : -1; } void process() { if (bus) while (sd_bus_process(bus, nullptr) > 0) {} } }; /* ── LEAdvertisement1 property handlers ────────────────────────────── */ static int adv_get_type(sd_bus*, const char*, const char*, const char*, sd_bus_message* reply, void*, sd_bus_error*) { return sd_bus_message_append(reply, "s", "peripheral"); } static int adv_get_uuids(sd_bus*, const char*, const char*, const char*, sd_bus_message* reply, void*, sd_bus_error*) { int r = sd_bus_message_open_container(reply, 'a', "s"); if (r < 0) return r; r = sd_bus_message_append(reply, "s", BLEBUS_SERVICE_UUID); if (r < 0) return r; return sd_bus_message_close_container(reply); } static int adv_get_name(sd_bus*, const char*, const char*, const char*, sd_bus_message* reply, void* ud, sd_bus_error*) { auto* adv = static_cast(ud); return sd_bus_message_append(reply, "s", adv->local_name.c_str()); } static int adv_release(sd_bus_message* msg, void*, sd_bus_error*) { return sd_bus_reply_method_return(msg, ""); } static const sd_bus_vtable adv_vtable[] = { SD_BUS_VTABLE_START(0), SD_BUS_PROPERTY("Type", "s", adv_get_type, 0, SD_BUS_VTABLE_PROPERTY_CONST), SD_BUS_PROPERTY("ServiceUUIDs", "as", adv_get_uuids, 0, SD_BUS_VTABLE_PROPERTY_CONST), SD_BUS_PROPERTY("LocalName", "s", adv_get_name, 0, SD_BUS_VTABLE_PROPERTY_CONST), SD_BUS_METHOD("Release", "", "", adv_release, 0), SD_BUS_VTABLE_END, }; constexpr const char* ADV_OBJ = "/org/blebus/ad0"; constexpr const char* ADV_IFACE = "org.bluez.LEAdvertisement1"; bool start_advertising(Advertiser& adv, std::string_view name) { adv.local_name = std::string(name); int r = sd_bus_open_system(&adv.bus); if (r < 0) return false; r = sd_bus_add_object_vtable(adv.bus, &adv.slot, ADV_OBJ, ADV_IFACE, adv_vtable, &adv); if (r < 0) { sd_bus_unref(adv.bus); adv.bus = nullptr; return false; } /* RegisterAdvertisement — BlueZ will query our object via GetAll * during sd_bus_call, which processes events internally. */ sd_bus_error error = SD_BUS_ERROR_NULL; sd_bus_message* m = nullptr; r = sd_bus_message_new_method_call( adv.bus, &m, "org.bluez", "/org/bluez/hci0", "org.bluez.LEAdvertisingManager1", "RegisterAdvertisement"); if (r < 0) return false; r = sd_bus_message_append(m, "o", ADV_OBJ); if (r >= 0) r = sd_bus_message_open_container(m, 'a', "{sv}"); if (r >= 0) r = sd_bus_message_close_container(m); if (r < 0) { sd_bus_message_unref(m); return false; } r = sd_bus_call(adv.bus, m, 5000000, &error, nullptr); /* 5s timeout */ sd_bus_message_unref(m); if (r < 0) { sd_bus_error_free(&error); return false; } adv.registered = true; return true; } void stop_advertising(Advertiser& adv) { if (!adv.bus) return; if (adv.registered) { sd_bus_error error = SD_BUS_ERROR_NULL; sd_bus_call_method(adv.bus, "org.bluez", "/org/bluez/hci0", "org.bluez.LEAdvertisingManager1", "UnregisterAdvertisement", &error, nullptr, "o", ADV_OBJ); sd_bus_error_free(&error); adv.registered = false; } if (adv.slot) { sd_bus_slot_unref(adv.slot); adv.slot = nullptr; } sd_bus_unref(adv.bus); adv.bus = nullptr; } /* ═══════════════════════════════════════════════════════════════════════ * Broker * ═══════════════════════════════════════════════════════════════════════ */ struct Broker { std::string name; int unix_listen_fd = -1; int ble_listen_fd = -1; std::thread thread; std::atomic running{false}; int pipe_fd[2] = {-1, -1}; size_t buf_size = 0; /* Client connections: unix socket + BLE L2CAP, tracked uniformly. */ int fds[MAX_READERS]; int count = 0; std::mutex mtx; Advertiser adv; Broker() { for (auto& fd : fds) fd = -1; } }; void broker_remove_client(Broker& b, int idx) { ::close(b.fds[idx]); b.count--; if (idx < b.count) b.fds[idx] = b.fds[b.count]; b.fds[b.count] = -1; } void broker_loop(Broker* b) { auto frame_buf = std::make_unique(b->buf_size); while (b->running.load()) { /* Build poll set: pipe + unix_listen + ble_listen + dbus + clients */ constexpr int FIXED_SLOTS = 4; /* pipe, unix, ble, dbus */ pollfd pfds[FIXED_SLOTS + MAX_READERS]; pfds[0] = {b->pipe_fd[0], POLLIN, 0}; pfds[1] = {b->unix_listen_fd, POLLIN, 0}; pfds[2] = {b->ble_listen_fd, static_cast(b->ble_listen_fd >= 0 ? POLLIN : 0), 0}; pfds[3] = {b->adv.fd(), static_cast(b->adv.fd() >= 0 ? POLLIN : 0), 0}; int nc; { std::lock_guard lock(b->mtx); nc = b->count; for (int i = 0; i < nc; i++) pfds[FIXED_SLOTS + i] = {b->fds[i], POLLIN, 0}; } int ret = ::poll(pfds, static_cast(FIXED_SLOTS + nc), 100); if (ret < 0) { if (errno == EINTR) continue; break; } if (ret == 0) continue; /* Shutdown signal */ if (pfds[0].revents & POLLIN) break; /* Process D-Bus events (advertising keepalive) */ if (pfds[3].revents & POLLIN) b->adv.process(); std::lock_guard lock(b->mtx); /* Accept new unix connections */ if (pfds[1].revents & POLLIN) { int cfd = ::accept4(b->unix_listen_fd, nullptr, nullptr, SOCK_NONBLOCK | SOCK_CLOEXEC); if (cfd >= 0) { if (b->count < static_cast(MAX_READERS)) b->fds[b->count++] = cfd; else ::close(cfd); } } /* Accept new BLE connections */ if (b->ble_listen_fd >= 0 && (pfds[2].revents & POLLIN)) { int cfd = ::accept4(b->ble_listen_fd, nullptr, nullptr, SOCK_NONBLOCK | SOCK_CLOEXEC); if (cfd >= 0) { if (b->count < static_cast(MAX_READERS)) b->fds[b->count++] = cfd; else ::close(cfd); } } /* Receive + broadcast (SEQPACKET: one recv = one message) */ bool dead[MAX_READERS] = {}; for (int i = 0; i < nc && i < b->count; i++) { if (!(pfds[FIXED_SLOTS + i].revents & (POLLIN | POLLHUP | POLLERR))) continue; ssize_t n = ::recv(b->fds[i], frame_buf.get(), b->buf_size, 0); if (n <= 0) { dead[i] = true; continue; } /* Broadcast to ALL clients (including sender) */ for (int j = 0; j < b->count; j++) { if (dead[j]) continue; ssize_t s = ::send(b->fds[j], frame_buf.get(), static_cast(n), MSG_DONTWAIT | MSG_NOSIGNAL); if (s < 0 && errno != EAGAIN && errno != EWOULDBLOCK) dead[j] = true; else if (s >= 0 && static_cast(s) < static_cast(n)) dead[j] = true; } } /* Remove dead clients (reverse order for stable indices) */ for (int i = b->count - 1; i >= 0; i--) { if (dead[i]) broker_remove_client(*b, i); } } /* Final cleanup: close all client sockets */ std::lock_guard lock(b->mtx); for (int i = 0; i < b->count; i++) ::close(b->fds[i]); b->count = 0; } /* ── Broker registry (process-global) ────────────────────────────────── */ std::unique_ptr g_brokers[MAX_BROKERS]; int g_broker_count = 0; std::mutex g_lock; } // anonymous namespace /* ═══════════════════════════════════════════════════════════════════════ */ /* PUBLIC API */ /* ═══════════════════════════════════════════════════════════════════════ */ void create(std::string_view name, size_t size) { if (name.empty() || name.size() >= MAX_NAME) throw std::system_error(EINVAL, std::system_category()); if (size == 0) size = DEFAULT_SIZE; std::lock_guard lock(g_lock); for (int i = 0; i < g_broker_count; i++) { if (g_brokers[i]->name == name) throw std::system_error(EEXIST, std::system_category()); } if (g_broker_count >= MAX_BROKERS) throw std::system_error(ENOMEM, std::system_category()); auto b = std::make_unique(); b->name = std::string(name); b->buf_size = size; /* ── Unix socket listener (local clients) ─────────────────────── */ std::string path = sock_path(name); ::unlink(path.c_str()); b->unix_listen_fd = ::socket(AF_UNIX, SOCK_SEQPACKET | SOCK_NONBLOCK | SOCK_CLOEXEC, 0); if (b->unix_listen_fd < 0) throw std::system_error(errno, std::system_category()); sockaddr_un uaddr{}; uaddr.sun_family = AF_UNIX; std::strncpy(uaddr.sun_path, path.c_str(), sizeof(uaddr.sun_path) - 1); if (::bind(b->unix_listen_fd, reinterpret_cast(&uaddr), sizeof(uaddr)) < 0) { int e = errno; ::close(b->unix_listen_fd); throw std::system_error(e, std::system_category()); } if (::listen(b->unix_listen_fd, static_cast(MAX_READERS)) < 0) { int e = errno; ::close(b->unix_listen_fd); ::unlink(path.c_str()); throw std::system_error(e, std::system_category()); } /* ── BLE L2CAP CoC listener (remote BLE clients) ──────────────── */ b->ble_listen_fd = ::socket(AF_BLUETOOTH, SOCK_SEQPACKET | SOCK_NONBLOCK | SOCK_CLOEXEC, BTPROTO_L2CAP); if (b->ble_listen_fd >= 0) { /* No pairing required — Antheos does its own auth (Z-verb). */ struct bt_security sec{}; sec.level = BT_SECURITY_LOW; ::setsockopt(b->ble_listen_fd, SOL_BLUETOOTH, BT_SECURITY, &sec, sizeof(sec)); /* Request maximum MTU for L2CAP CoC SDU size. */ uint16_t mtu = 65535; ::setsockopt(b->ble_listen_fd, SOL_BLUETOOTH, BT_RCVMTU, &mtu, sizeof(mtu)); ::setsockopt(b->ble_listen_fd, SOL_BLUETOOTH, BT_SNDMTU, &mtu, sizeof(mtu)); sockaddr_l2 laddr{}; laddr.l2_family = AF_BLUETOOTH; laddr.l2_psm = htobs(PSM); laddr.l2_bdaddr = {}; /* BDADDR_ANY */ laddr.l2_cid = 0; laddr.l2_bdaddr_type = BDADDR_LE_PUBLIC; if (::bind(b->ble_listen_fd, reinterpret_cast(&laddr), sizeof(laddr)) < 0 || ::listen(b->ble_listen_fd, static_cast(MAX_READERS)) < 0) { ::close(b->ble_listen_fd); b->ble_listen_fd = -1; /* BLE unavailable — unix-only mode (still useful for testing). */ } } /* ── BLE advertising ──────────────────────────────────────────── */ if (b->ble_listen_fd >= 0) start_advertising(b->adv, name); /* Advertising failure is non-fatal — BLE clients can still connect * if they know the address. Unix socket always works. */ /* ── Shutdown pipe + broker thread ────────────────────────────── */ if (::pipe2(b->pipe_fd, O_CLOEXEC) < 0) { int e = errno; ::close(b->unix_listen_fd); ::unlink(path.c_str()); if (b->ble_listen_fd >= 0) ::close(b->ble_listen_fd); stop_advertising(b->adv); throw std::system_error(e, std::system_category()); } b->running.store(true); Broker* bp = b.get(); b->thread = std::thread(broker_loop, bp); g_brokers[g_broker_count++] = std::move(b); } void destroy(std::string_view name) { if (name.empty()) throw std::system_error(EINVAL, std::system_category()); std::unique_lock lock(g_lock); for (int i = 0; i < g_broker_count; i++) { if (g_brokers[i]->name != name) continue; auto& b = g_brokers[i]; b->running.store(false); (void)::write(b->pipe_fd[1], "x", 1); lock.unlock(); if (b->thread.joinable()) b->thread.join(); lock.lock(); stop_advertising(b->adv); ::close(b->unix_listen_fd); if (b->ble_listen_fd >= 0) ::close(b->ble_listen_fd); ::close(b->pipe_fd[0]); ::close(b->pipe_fd[1]); ::unlink(sock_path(b->name).c_str()); g_broker_count--; if (i < g_broker_count) g_brokers[i] = std::move(g_brokers[g_broker_count]); g_brokers[g_broker_count].reset(); return; } throw std::system_error(ENOENT, std::system_category()); } /* ── Bus::Impl ───────────────────────────────────────────────────────── */ struct Bus::Impl { std::string bus_name; int fd = -1; char reader_name[READER_NAME_LEN] = {}; std::mutex mtx; std::unique_ptr recv_buf; size_t recv_cap = 0; ~Impl() { if (fd >= 0) ::close(fd); } }; Bus::Bus(std::string_view name) : impl_(std::make_unique()) { if (name.empty()) throw std::system_error(EINVAL, std::system_category()); std::string path = sock_path(name); int fd = ::socket(AF_UNIX, SOCK_SEQPACKET | SOCK_CLOEXEC, 0); if (fd < 0) throw std::system_error(errno, std::system_category()); sockaddr_un addr{}; addr.sun_family = AF_UNIX; std::strncpy(addr.sun_path, path.c_str(), sizeof(addr.sun_path) - 1); if (::connect(fd, reinterpret_cast(&addr), sizeof(addr)) < 0) { int e = errno; ::close(fd); throw std::system_error(e, std::system_category()); } impl_->bus_name = std::string(name); impl_->fd = fd; impl_->recv_cap = DEFAULT_SIZE; impl_->recv_buf = std::make_unique(impl_->recv_cap); } Bus::~Bus() = default; Bus::Bus(Bus&& o) noexcept = default; Bus& Bus::operator=(Bus&& o) noexcept = default; size_t Bus::write(const uint8_t* data, size_t len) { if (!data) throw std::system_error(EINVAL, std::system_category()); if (len > MAX_MSG_SIZE) throw std::system_error(EMSGSIZE, std::system_category()); std::lock_guard lock(impl_->mtx); ssize_t ret = send_frame(impl_->fd, data, len); if (ret < 0) throw std::system_error(errno, std::system_category()); return static_cast(ret); } size_t Bus::write(std::string_view sv) { return write(reinterpret_cast(sv.data()), sv.size()); } size_t Bus::read(uint8_t* buf, size_t len) { if (!buf) throw std::system_error(EINVAL, std::system_category()); std::lock_guard lock(impl_->mtx); /* SEQPACKET: one recv = one complete message (or EAGAIN). */ ssize_t n = ::recv(impl_->fd, impl_->recv_buf.get(), impl_->recv_cap, MSG_DONTWAIT); if (n <= 0) return 0; const uint8_t* payload; size_t plen = extract_payload(impl_->recv_buf.get(), static_cast(n), &payload); if (plen == 0) return 0; size_t copy = std::min(plen, len); std::memcpy(buf, payload, copy); return copy; } size_t Bus::read_wait(uint8_t* buf, size_t len, int timeout_ms) { if (!buf) throw std::system_error(EINVAL, std::system_category()); if (timeout_ms == 0) return read(buf, len); timespec deadline{}; if (timeout_ms > 0) { clock_gettime(CLOCK_MONOTONIC, &deadline); deadline.tv_sec += timeout_ms / 1000; deadline.tv_nsec += static_cast(timeout_ms % 1000) * 1000000L; if (deadline.tv_nsec >= 1000000000L) { deadline.tv_sec++; deadline.tv_nsec -= 1000000000L; } } for (;;) { { std::lock_guard lock(impl_->mtx); ssize_t n = ::recv(impl_->fd, impl_->recv_buf.get(), impl_->recv_cap, MSG_DONTWAIT); if (n > 0) { const uint8_t* payload; size_t plen = extract_payload(impl_->recv_buf.get(), static_cast(n), &payload); if (plen > 0) { size_t copy = std::min(plen, len); std::memcpy(buf, payload, copy); return copy; } } } int poll_ms = -1; if (timeout_ms > 0) { timespec now{}; clock_gettime(CLOCK_MONOTONIC, &now); long rem = (deadline.tv_sec - now.tv_sec) * 1000 + (deadline.tv_nsec - now.tv_nsec) / 1000000; if (rem <= 0) return 0; poll_ms = static_cast(rem); } pollfd pfd = {impl_->fd, POLLIN, 0}; int ret = ::poll(&pfd, 1, poll_ms); if (ret < 0) { if (errno == EINTR) continue; throw std::system_error(errno, std::system_category()); } if (ret == 0) return 0; } } void Bus::set_reader_name(std::string_view label) { size_t copy = std::min(label.size(), READER_NAME_LEN - 1); std::memcpy(impl_->reader_name, label.data(), copy); impl_->reader_name[copy] = '\0'; } std::string_view Bus::name() const { return impl_->bus_name; } } // namespace blebus