Skip to main content

lego_powered_up/iodevice/
sensor.rs

1//! Support for generic sensors. Can be used with simple sensors
2//! like hub temp, voltage etc., or other devices without
3//! higher level support.
4
5use async_trait::async_trait;
6use core::fmt::Debug;
7use tokio::sync::broadcast;
8use tokio::task::JoinHandle;
9
10use super::Basic;
11use crate::device_trait;
12use crate::error::{Error, Result};
13use crate::notifications::DatasetType;
14
15device_trait!(GenericSensor, [
16    fn check_dataset(&self, mode: u8, datasettype: DatasetType) -> Result<()>;,
17
18    async fn enable_8bit_sensor(
19        &self,
20        mode: u8,
21        delta: u32,
22    ) -> Result<(broadcast::Receiver<Vec<i8>>, JoinHandle<()>)> {
23        match self.check_dataset(mode, DatasetType::Bits8) {
24            Ok(()) => (),
25            _ => {
26                return Err(Error::NoneError(String::from(
27                    "Not an 8-bit sensor mode",
28                )))
29            }
30        }
31
32        self.device_mode(mode, delta, true).await?;
33
34        // Set up channel
35        let port_id = self.port();
36        let (tx, rx) = broadcast::channel::<Vec<i8>>(64);
37        match self.get_rx() {
38            Ok(mut rx_from_main) => {
39                let task = tokio::spawn(async move {
40                    while let Ok(msg) = rx_from_main.recv().await {
41                        if msg.port_id != port_id {
42                            continue;
43                        }
44                        let _ = tx.send(msg.data);
45                    }
46                });
47
48                Ok((rx, task))
49            }
50            _ => {
51                Err(Error::NoneError(String::from("No sender in device cache")))
52            }
53        }
54    },
55
56    async fn enable_16bit_sensor(
57        &self,
58        mode: u8,
59        delta: u32,
60    ) -> Result<(broadcast::Receiver<Vec<i16>>, JoinHandle<()>)> {
61        match self.check_dataset(mode, DatasetType::Bits16) {
62            Ok(()) => (),
63            _ => {
64                return Err(Error::NoneError(String::from(
65                    "Not a 16-bit sensor mode",
66                )))
67            }
68        }
69        self.device_mode(mode, delta, true).await?;
70
71        // Set up channel
72        let port_id = self.port();
73        let (tx, rx) = broadcast::channel::<Vec<i16>>(64);
74        match self.get_rx() {
75            Ok(mut rx_from_main) => {
76                let task = tokio::spawn(async move {
77                    while let Ok(data) = rx_from_main.recv().await {
78                        if data.port_id != port_id {
79                            continue;
80                        }
81                        let mut converted: Vec<i16> = Vec::new();
82
83
84                        let cycles = &data.data.len() / 2;
85                        let mut it = data.data.into_iter();
86                        for _ in 0..cycles {
87                            converted.push(i16::from_le_bytes([
88                                it.next().unwrap() as u8,
89                                it.next().unwrap() as u8,
90                            ]));
91                        }
92
93                        let _ = tx.send(converted);
94                    }
95                });
96
97                Ok((rx, task))
98            }
99            _ => {
100                Err(Error::NoneError(String::from("No sender in device cache")))
101            }
102        }
103    },
104
105    async fn enable_32bit_sensor(
106        &self,
107        mode: u8,
108        delta: u32,
109    ) -> Result<(broadcast::Receiver<Vec<i32>>, JoinHandle<()>)> {
110        match self.check_dataset(mode, DatasetType::Bits32) {
111            Ok(()) => (),
112            _ => {
113                return Err(Error::NoneError(String::from(
114                    "Not a 32-bit sensor mode",
115                )))
116            }
117        }
118        self.device_mode(mode, delta, true).await?;
119
120        // Set up channel
121        let port_id = self.port();
122        let (tx, rx) = broadcast::channel::<Vec<i32>>(64);
123        match self.get_rx() {
124            Ok(mut rx_from_main) => {
125                let task = tokio::spawn(async move {
126                    while let Ok(data) = rx_from_main.recv().await {
127                        if data.port_id != port_id {
128                            continue;
129                        }
130                        // println!("32bit data: {:?}", &data.data);
131
132                        let mut converted: Vec<i32> = Vec::new();
133
134                        let cycles = &data.data.len() / 4;
135                        let mut it = data.data.into_iter();
136                        for _ in 0..cycles {
137                            converted.push(i32::from_le_bytes([
138                                it.next().unwrap() as u8,
139                                it.next().unwrap() as u8,
140                                it.next().unwrap() as u8,
141                                it.next().unwrap() as u8,
142                            ]));
143                        }
144
145                        let _ = tx.send(converted);
146                    }
147                });
148
149                Ok((rx, task))
150            }
151            _ => {
152                Err(Error::NoneError(String::from("No sender in device cache")))
153            }
154        }
155    },
156
157    fn raw_channel(
158        &self,
159    ) -> Result<(broadcast::Receiver<Vec<i8>>, JoinHandle<()>)> {
160        let port_id = self.port();
161        let (tx, rx) = broadcast::channel::<Vec<i8>>(64);
162        match self.get_rx() {
163            Ok(mut rx_from_main) => {
164                let task = tokio::spawn(async move {
165                    while let Ok(msg) = rx_from_main.recv().await {
166                        if msg.port_id != port_id {
167                            continue;
168                        }
169                        let _ = tx.send(msg.data);
170                    }
171                });
172
173                Ok((rx, task))
174            }
175            _ => {
176                Err(Error::NoneError(String::from("No sender in device cache")))
177            }
178        }
179    }
180
181]);