lego_powered_up/iodevice/
sensor.rs1use 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 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 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 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 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]);