| 34 | type Peripheral = Peripheral; |
| 35 | |
| 36 | async fn events(&self) -> Result<Pin<Box<dyn Stream<Item = CentralEvent> + Send>>> { |
| 37 | // There's a race between getting this event stream and getting the current set of devices. |
| 38 | // Get the stream first, on the basis that it's better to have a duplicate DeviceDiscovered |
| 39 | // event than to miss one. It's unlikely to happen in any case. |
| 40 | let events = self.session.adapter_event_stream(&self.adapter).await?; |
| 41 | |
| 42 | // Synthesise `DeviceDiscovered' and `DeviceConnected` events for existing peripherals. |
| 43 | let devices = self.session.get_devices().await?; |
| 44 | let adapter_id = self.adapter.clone(); |
| 45 | let initial_events = stream::iter( |
| 46 | devices |
| 47 | .into_iter() |
| 48 | .filter(move |device| device.id.adapter() == adapter_id) |
| 49 | .flat_map(|device| { |
| 50 | let peripheral_id: PeripheralId = device.id.into(); |
| 51 | let mut events = vec![CentralEvent::DeviceDiscovered(peripheral_id.clone())]; |
| 52 | if !device.services.is_empty() { |
| 53 | events.push(CentralEvent::ServicesAdvertisement { |
| 54 | id: peripheral_id.clone(), |
| 55 | services: device.services, |
| 56 | }); |
| 57 | } |
| 58 | if device.connected { |
| 59 | events.push(CentralEvent::DeviceConnected(peripheral_id)); |
| 60 | } |
| 61 | events.into_iter() |
| 62 | }), |
| 63 | ); |
| 64 | |
| 65 | let session = self.session.clone(); |
| 66 | let adapter_id = self.adapter.clone(); |
| 67 | let events = events |
| 68 | .filter_map(move |event| central_events(event, session.clone(), adapter_id.clone())) |
| 69 | .flat_map(stream::iter); |
| 70 | |
| 71 | Ok(Box::pin(initial_events.chain(events))) |
| 72 | } |
| 73 | |
| 74 | async fn start_scan(&self, filter: ScanFilter) -> Result<()> { |
| 75 | let filter = DiscoveryFilter { |