Skip to content

Commit 7751f46

Browse files
Merge pull request #15 from tyler-potyondy/fix/rpc-send-recv
feat: add polling infastructure for events and callbacks from server
2 parents 33768e2 + 8e2150d commit 7751f46

9 files changed

Lines changed: 812 additions & 546 deletions

File tree

src/ble/mod.rs

Lines changed: 371 additions & 496 deletions
Large diffs are not rendered by default.

src/lib.rs

Lines changed: 285 additions & 29 deletions
Large diffs are not rendered by default.

src/transport.rs

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,7 @@ pub trait AsyncTransport {
4848
/// Write bytes to the transport
4949
///
5050
/// Should block until all bytes are written or an error occurs.
51-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error>;
51+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error>;
5252

5353
/// Read bytes from the transport into the provided buffer
5454
///
@@ -58,6 +58,31 @@ pub trait AsyncTransport {
5858

5959
/// Delay for the specified number of milliseconds
6060
async fn delay_ms(&mut self, ms: u32);
61+
62+
/// Returns `true` if bytes are already buffered and [`read`](Self::read)
63+
/// can make progress without suspending.
64+
///
65+
/// **Implementations must be backed by a persistent ring buffer** (e.g. a
66+
/// DMA-filled circular buffer or interrupt-driven FIFO). The check is
67+
/// synchronous and must not consume any bytes.
68+
///
69+
/// This is used by [`RpcClient::has_data`] to let callers in a two-task
70+
/// embassy setup skip taking a mutex when there is nothing to process:
71+
///
72+
/// ```ignore
73+
/// loop {
74+
/// let mut ble = ble_mutex.lock().await;
75+
/// if ble.has_data() {
76+
/// let evt = ble.next_event().await?; // completes quickly
77+
/// drop(ble);
78+
/// handle(evt);
79+
/// } else {
80+
/// drop(ble); // release immediately — nothing to do
81+
/// embassy_futures::yield_now().await;
82+
/// }
83+
/// }
84+
/// ```
85+
fn has_buffered_data(&mut self) -> bool;
6186
}
6287

6388
#[derive(Debug)]

src/uart_transport.rs

Lines changed: 31 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -105,10 +105,33 @@ pub trait Uart {
105105
async fn read(&mut self, buf: &mut [u8]) -> Result<usize, Self::Error>;
106106

107107
/// Write raw bytes to the UART.
108-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error>;
108+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error>;
109109

110110
/// Delay for the given number of milliseconds.
111111
async fn delay_ms(&mut self, ms: u32);
112+
113+
/// Returns `true` if at least one byte is already sitting in the
114+
/// hardware/DMA ring buffer and can be returned by [`read`](Self::read)
115+
/// without suspending.
116+
///
117+
/// Implementations **must** be backed by a persistent ring buffer (e.g. a
118+
/// DMA-filled circular buffer or an interrupt-driven FIFO). The check must
119+
/// be synchronous and must not consume any bytes.
120+
///
121+
/// For Embassy targets, delegate to [`embedded_io::ReadReady`] which checks
122+
/// the DMA ring buffer's read/write pointers without consuming bytes:
123+
///
124+
/// ```ignore
125+
/// fn has_buffered_data(&mut self) -> bool {
126+
/// use embedded_io::ReadReady;
127+
/// self.uart.read_ready().unwrap_or(false)
128+
/// }
129+
/// ```
130+
///
131+
/// If `ReadReady` is not available for your peripheral, return `false` as a
132+
/// safe conservative fallback — the two-task pattern remains correct, it
133+
/// just polls via `yield_now()` on every loop iteration.
134+
fn has_buffered_data(&mut self) -> bool;
112135
}
113136

114137
/// UART transport wrapper that implements [`AsyncTransport`] on top of a [`Uart`].
@@ -139,10 +162,16 @@ impl<Inner: Uart> AsyncTransport for UartTransport<Inner> {
139162
type TxTransportPacket<'a> = UartTxTransport<'a>;
140163
type RxTransportPacket<'a> = UartRxTransport<'a>;
141164

142-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error> {
165+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error> {
143166
self.inner.write(data).await
144167
}
145168

169+
/// Returns `true` if a complete HDLC frame is already accumulated in the
170+
/// internal buffer, or if the underlying [`Uart`] has bytes ready.
171+
fn has_buffered_data(&mut self) -> bool {
172+
hdlc_frame_complete(&self.rx_buf[..self.rx_len]) || self.inner.has_buffered_data()
173+
}
174+
146175
/// Accumulate raw bytes from the inner transport until a complete HDLC frame
147176
/// is present, then deliver all bytes up through the closing delimiter.
148177
///

tests/batched_packet_test.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -144,7 +144,7 @@ impl DelayedOneShotTransport {
144144
impl Uart for DelayedOneShotTransport {
145145
type Error = MockError;
146146

147-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error> {
147+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error> {
148148
self.writes.lock().unwrap().extend_from_slice(data);
149149
Ok(data.len())
150150
}
@@ -164,6 +164,8 @@ impl Uart for DelayedOneShotTransport {
164164
}
165165

166166
async fn delay_ms(&mut self, _ms: u32) {}
167+
168+
fn has_buffered_data(&mut self) -> bool { false }
167169
}
168170

169171
// ── tests ─────────────────────────────────────────────────────────────────────

tests/ble_bt_enable_test.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ impl DummyUart {
4545
impl Uart for DummyUart {
4646
type Error = DummyError;
4747

48-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error> {
48+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error> {
4949
let mut state = self.state.lock().unwrap();
5050
state.sent_packets.push(data.to_vec());
5151
Ok(data.len())
@@ -66,6 +66,8 @@ impl Uart for DummyUart {
6666
}
6767

6868
async fn delay_ms(&mut self, _ms: u32) {}
69+
70+
fn has_buffered_data(&mut self) -> bool { false }
6971
}
7072

7173
/// Minimal async executor for this test - same pattern as in integration_test.rs.

tests/integration_test.rs

Lines changed: 80 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -12,9 +12,9 @@ use nrf_rpc::ble::cgm::{
1212
CgmMeasurement, encode_uuid_16,
1313
};
1414
use nrf_rpc::ble::{
15-
BT_GATT_CCC_NOTIFY, BT_LE_SCAN_TYPE_ACTIVE, BtConnLeCreateParam, BtGattDiscoverParams,
16-
BtGattDiscoverType, BtGattSubscribeParams, BtLeConnParam, BtLeScanParam, GattDiscoverResult,
17-
ScanResultData,
15+
BleEvent, BT_GATT_CCC_NOTIFY, BT_LE_SCAN_TYPE_ACTIVE, BtConnLeCreateParam,
16+
BtGattDiscoverParams, BtGattDiscoverType, BtGattSubscribeParams, BtLeConnParam, BtLeScanParam,
17+
GattDiscoverResult, ScanResultData,
1818
};
1919
use babble_bridge::{LogOutput, TestProcesses, spawn_zephyr_rpc_server_with_socat};
2020
use nrf_rpc::{RpcClient, TransportError, ble::Ble, uart_transport::{Uart, UartTransport}};
@@ -191,7 +191,7 @@ impl Uart for MockUart {
191191
}
192192
}
193193

194-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error> {
194+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error> {
195195
println!(
196196
"MockUart: Sending {} bytes to {}: {:02X?}",
197197
data.len(),
@@ -226,6 +226,10 @@ impl Uart for MockUart {
226226
async fn delay_ms(&mut self, ms: u32) {
227227
std::thread::sleep(Duration::from_millis(ms as u64));
228228
}
229+
230+
fn has_buffered_data(&mut self) -> bool {
231+
!self.rx_buffer.lock().unwrap().is_empty()
232+
}
229233
}
230234

231235
/// Helper to convert hex string to bytes
@@ -742,7 +746,16 @@ fn test_cgm_full_central_flow_inner() {
742746
let max_scan_results = 50;
743747

744748
for i in 0..max_scan_results {
745-
let result = embassy_futures::block_on(ble.wait_for_scan_result());
749+
let result = embassy_futures::block_on(async {
750+
loop {
751+
match ble.next_event().await {
752+
Ok(BleEvent::ScanResult(s)) => return Ok(s),
753+
Ok(BleEvent::ScanTimeout) => return Err(nrf_rpc::ble::BleError::RpcError),
754+
Ok(_) => continue,
755+
Err(e) => return Err(e),
756+
}
757+
}
758+
});
746759
match result {
747760
Ok(scan) => {
748761
let name = scan.device_name().unwrap_or("<unknown>");
@@ -809,7 +822,16 @@ fn test_cgm_full_central_flow_inner() {
809822
// Step 6: Wait for the "connected" callback event
810823
// ------------------------------------------------------------------
811824
println!("[Step 6] Waiting for connection event...");
812-
let conn_event = embassy_futures::block_on(ble.wait_for_connection());
825+
let conn_event = embassy_futures::block_on(async {
826+
for _ in 0..20 {
827+
match ble.next_event().await {
828+
Ok(BleEvent::Connected(e)) => return Ok(e),
829+
Ok(_) => continue,
830+
Err(e) => return Err(e),
831+
}
832+
}
833+
Err(nrf_rpc::ble::BleError::RpcError)
834+
});
813835
assert!(
814836
conn_event.is_ok(),
815837
"Did not receive connection event: {:?}",
@@ -846,7 +868,22 @@ fn test_cgm_full_central_flow_inner() {
846868

847869
// Wait for the SMP exchange to complete (passkey exchange + security level 4)
848870
println!("[Step 6a] Waiting for security level 4...");
849-
let result = embassy_futures::block_on(ble.wait_for_security_level(4));
871+
let result = embassy_futures::block_on(async {
872+
for _ in 0..30 {
873+
match ble.next_event().await {
874+
Ok(BleEvent::SecurityChanged { level, err: 0 }) if level >= 4 => return Ok(level),
875+
Ok(BleEvent::PasskeyConfirm(_)) => {
876+
let _ = ble.bt_conn_auth_passkey_confirm().await;
877+
}
878+
Ok(BleEvent::PairingConfirm) => {
879+
let _ = ble.bt_conn_auth_pairing_confirm().await;
880+
}
881+
Ok(_) => continue,
882+
Err(e) => return Err(e),
883+
}
884+
}
885+
Err(nrf_rpc::ble::BleError::RpcError)
886+
});
850887
assert!(
851888
result.is_ok(),
852889
"Failed to achieve security level 4: {:?}",
@@ -884,7 +921,15 @@ fn test_cgm_full_central_flow_inner() {
884921
let mut cgm_service_end_handle: Option<u16> = None;
885922

886923
for i in 0..20 {
887-
let result = embassy_futures::block_on(ble.wait_for_gatt_discover_result());
924+
let result = embassy_futures::block_on(async {
925+
loop {
926+
match ble.next_event().await {
927+
Ok(BleEvent::GattDiscovery(r)) => return Ok(r),
928+
Ok(_) => continue,
929+
Err(e) => return Err(e),
930+
}
931+
}
932+
});
888933
match result {
889934
Ok(GattDiscoverResult::Service {
890935
handle,
@@ -951,7 +996,15 @@ fn test_cgm_full_central_flow_inner() {
951996
let mut cgm_meas_value_handle: Option<u16> = None;
952997

953998
for i in 0..20 {
954-
let result = embassy_futures::block_on(ble.wait_for_gatt_discover_result());
999+
let result = embassy_futures::block_on(async {
1000+
loop {
1001+
match ble.next_event().await {
1002+
Ok(BleEvent::GattDiscovery(r)) => return Ok(r),
1003+
Ok(_) => continue,
1004+
Err(e) => return Err(e),
1005+
}
1006+
}
1007+
});
9551008
match result {
9561009
Ok(GattDiscoverResult::Characteristic {
9571010
handle,
@@ -1017,7 +1070,15 @@ fn test_cgm_full_central_flow_inner() {
10171070
let mut ccc_handle: Option<u16> = None;
10181071

10191072
for i in 0..20 {
1020-
let result = embassy_futures::block_on(ble.wait_for_gatt_discover_result());
1073+
let result = embassy_futures::block_on(async {
1074+
loop {
1075+
match ble.next_event().await {
1076+
Ok(BleEvent::GattDiscovery(r)) => return Ok(r),
1077+
Ok(_) => continue,
1078+
Err(e) => return Err(e),
1079+
}
1080+
}
1081+
});
10211082
match result {
10221083
Ok(GattDiscoverResult::Descriptor { handle, uuid_16 }) => {
10231084
println!(" Desc #{}: UUID=0x{:04X} handle={}", i, uuid_16, handle);
@@ -1085,7 +1146,15 @@ fn test_cgm_full_central_flow_inner() {
10851146
let max_notification_attempts = 30;
10861147

10871148
for attempt in 0..max_notification_attempts {
1088-
let result = embassy_futures::block_on(ble.wait_for_gatt_notification());
1149+
let result = embassy_futures::block_on(async {
1150+
loop {
1151+
match ble.next_event().await {
1152+
Ok(BleEvent::GattNotification(n)) => return Ok(n),
1153+
Ok(_) => continue,
1154+
Err(e) => return Err(e),
1155+
}
1156+
}
1157+
});
10891158
match result {
10901159
Ok(notif) => {
10911160
println!(

tests/transport_retry_test.rs

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -94,13 +94,14 @@ struct AlwaysErrTransport;
9494

9595
impl Uart for AlwaysErrTransport {
9696
type Error = MockError;
97-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error> {
97+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error> {
9898
Ok(data.len())
9999
}
100100
async fn read(&mut self, _buf: &mut [u8]) -> Result<usize, Self::Error> {
101101
Err(MockError)
102102
}
103103
async fn delay_ms(&mut self, _ms: u32) {}
104+
fn has_buffered_data(&mut self) -> bool { false }
104105
}
105106

106107
/// `read()` returns `Err` for the first `fail_count` calls, then returns a
@@ -124,7 +125,7 @@ impl FailThenSucceedTransport {
124125

125126
impl Uart for FailThenSucceedTransport {
126127
type Error = MockError;
127-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error> {
128+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error> {
128129
Ok(data.len())
129130
}
130131
async fn read(&mut self, buf: &mut [u8]) -> Result<usize, Self::Error> {
@@ -138,6 +139,7 @@ impl Uart for FailThenSucceedTransport {
138139
Ok(n)
139140
}
140141
async fn delay_ms(&mut self, _ms: u32) {}
142+
fn has_buffered_data(&mut self) -> bool { false }
141143
}
142144

143145
/// `read()` always returns `Ok(0)` — no data available.
@@ -150,13 +152,14 @@ struct NoDataTransport;
150152

151153
impl Uart for NoDataTransport {
152154
type Error = MockError;
153-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error> {
155+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error> {
154156
Ok(data.len())
155157
}
156158
async fn read(&mut self, _buf: &mut [u8]) -> Result<usize, Self::Error> {
157159
Ok(0)
158160
}
159161
async fn delay_ms(&mut self, _ms: u32) {}
162+
fn has_buffered_data(&mut self) -> bool { false }
160163
}
161164

162165
/// `read()` returns `Ok(0)` for the first `empty_count` calls, then returns a
@@ -180,7 +183,7 @@ impl EmptyThenSucceedTransport {
180183

181184
impl Uart for EmptyThenSucceedTransport {
182185
type Error = MockError;
183-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error> {
186+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error> {
184187
Ok(data.len())
185188
}
186189
async fn read(&mut self, buf: &mut [u8]) -> Result<usize, Self::Error> {
@@ -194,6 +197,7 @@ impl Uart for EmptyThenSucceedTransport {
194197
Ok(n)
195198
}
196199
async fn delay_ms(&mut self, _ms: u32) {}
200+
fn has_buffered_data(&mut self) -> bool { false }
197201
}
198202

199203
// ── tests ─────────────────────────────────────────────────────────────────────

tests/utils/mod.rs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,10 +53,14 @@ impl Uart for MockUart {
5353
Ok(0)
5454
}
5555

56-
async fn write(&mut self, data: &mut [u8]) -> Result<usize, Self::Error> {
56+
async fn write(&mut self, data: &[u8]) -> Result<usize, Self::Error> {
5757
self.transmitted.borrow_mut().extend_from_slice(data);
5858
Ok(data.len())
5959
}
6060

6161
async fn delay_ms(&mut self, _ms: u32) {}
62+
63+
fn has_buffered_data(&mut self) -> bool {
64+
false
65+
}
6266
}

0 commit comments

Comments
 (0)