crazyflie_lib/subsystems/memory/
mod.rs1use crate::{crtp_utils::WaitForPacket, Error, Result};
12use crazyflie_link::Packet;
13use flume as channel;
14use std::{collections::HashMap, convert::{TryFrom, TryInto}};
15use std::sync::Arc;
16use tokio::sync::Mutex;
17
18mod memory_types;
19mod eeprom_config;
20mod deckmem;
21mod raw;
22mod ow;
23mod trajectory;
24mod lighthouse;
25mod loco2;
26mod led_driver;
27
28use crate::crazyflie::MEMORY_PORT;
29
30pub use memory_types::*;
31pub use eeprom_config::*;
32pub use deckmem::*;
33pub use raw::*;
34pub use ow::*;
35pub use trajectory::*;
36pub use lighthouse::*;
37pub use loco2::*;
38pub use led_driver::*;
39
40#[derive(Debug)]
45pub struct Memory {
46 memories: Vec<MemoryDevice>,
47 backends: Vec<Mutex<Option<MemoryBackend>>>,
48 memory_read_dispatcher: MemoryDispatcher,
49 memory_write_dispatcher: MemoryDispatcher,
50}
51
52const INFO_CHANNEL: u8 = 0;
53const READ_CHANNEL: u8 = 1;
54const WRITE_CHANNEL: u8 = 2;
55
56const _CMD_INFO_VER: u8 = 0;
57const CMD_INFO_NBR: u8 = 1;
58const CMD_INFO_DETAILS: u8 = 2;
59
60#[derive(Debug)]
61struct MemoryDispatcher {
62 senders: Arc<Mutex<HashMap<u8, channel::Sender<Packet>>>>,
63}
64
65impl MemoryDispatcher {
66 fn new(downlink: channel::Receiver<Packet>, channel: u8) -> Self {
67
68 let senders: Arc<Mutex<HashMap<u8, channel::Sender<Packet>>>> = Arc::new(Mutex::new(HashMap::new()));
69 let internal_senders = senders.clone();
70
71 tokio::spawn(async move {
72 while let Ok(pk) = downlink.recv_async().await {
73 if pk.get_channel() == channel {
74 let memory_id = pk.get_data()[0];
75 if let Some(sender) = internal_senders.lock().await.get(&memory_id) {
76 let _ = sender.send_async(pk).await;
77 } else {
78 println!("Error: Received memory read response for unknown memory ID {}", memory_id);
79 break;
80 }
81 } else {
82 println!("Error: Received packet on unexpected channel {}", pk.get_channel());
83 break;
84 }
85 }
86 internal_senders.lock().await.clear();
87 });
88
89 Self {
90 senders: senders,
91 }
92 }
93
94 async fn get_channel(&mut self, memory_id: u8) -> channel::Receiver<Packet> {
95 if !self.senders.lock().await.contains_key(&memory_id) {
96 let (tx, rx) = channel::unbounded();
97 self.senders.lock().await.insert(memory_id, tx);
98 rx
99 } else {
100 panic!("Channel for memory ID {} already exists", memory_id)
101 }
102 }
103}
104
105impl Memory {
106 pub(crate) async fn new(
107 downlink: channel::Receiver<Packet>,
108 uplink: channel::Sender<Packet>,
109 ) -> Result<Self> {
110 let (info_channel_downlink, read_channel_downlink, write_channel_downlink, _misc_downlink) =
111 crate::crtp_utils::crtp_channel_dispatcher(downlink);
112
113 let mut memory = Self {
114 memories: Vec::new(),
115 backends: Vec::new(),
116 memory_read_dispatcher: MemoryDispatcher::new(read_channel_downlink.clone(), READ_CHANNEL),
117 memory_write_dispatcher: MemoryDispatcher::new(write_channel_downlink.clone(), WRITE_CHANNEL),
118 };
119
120 memory.update_memories(uplink.clone(), info_channel_downlink).await?;
121
122 Ok(memory)
123 }
124
125 async fn update_memories(&mut self, uplink: channel::Sender<Packet>, downlink: channel::Receiver<Packet>) -> Result<()> {
126 let pk = Packet::new(MEMORY_PORT, INFO_CHANNEL, vec![CMD_INFO_NBR]);
127 uplink
128 .send_async(pk)
129 .await
130 .map_err(|_| Error::Disconnected)?;
131
132 let pk = downlink.wait_packet(MEMORY_PORT, INFO_CHANNEL, &[CMD_INFO_NBR]).await?;
133 let memory_count = pk.get_data()[1];
134
135 for i in 0..memory_count {
136 let pk = Packet::new(MEMORY_PORT, INFO_CHANNEL, vec![CMD_INFO_DETAILS, i]);
137 uplink
138 .send_async(pk)
139 .await
140 .map_err(|_| Error::Disconnected)?;
141
142 let pk = downlink.wait_packet(MEMORY_PORT, INFO_CHANNEL, &[CMD_INFO_DETAILS, i]).await?;
143 let data = pk.get_data();
144 let memory_id = data[1];
145 let memory_type = MemoryType::try_from(data[2])?;
146 let memory_size = u32::from_le_bytes(data[3..7].try_into()?);
147 let raw_memory_serial = Vec::from(&data[7..]);
148
149 let memory_serial = if raw_memory_serial.iter().all(|&b| b == 0) {
150 None
151 } else {
152 Some(raw_memory_serial)
153 };
154
155 self.memories.push(MemoryDevice {
156 memory_id: memory_id,
157 memory_type: memory_type,
158 size: memory_size,
159 serial: memory_serial,
160 });
161
162 self.backends.push(Mutex::new(Some(MemoryBackend {
163 memory_id: memory_id,
164 memory_type: memory_type,
165 uplink: uplink.clone(),
166 read_downlink: self.memory_read_dispatcher.get_channel(memory_id).await,
167 write_downlink: self.memory_write_dispatcher.get_channel(memory_id).await,
168 })));
169 }
170 Ok(())
171
172 }
173
174 pub fn get_memories(&self, memory_type: Option<MemoryType>) -> Vec<&MemoryDevice> {
215 match memory_type {
216 Some(ty) => self.memories.iter().filter(|m| m.memory_type == ty).collect(),
217 None => self.memories.iter().collect(),
218 }
219 }
220
221 pub async fn open_memory<T: FromMemoryBackend>(&self, memory: MemoryDevice) -> Option<Result<T>> {
228 let backend = self.backends.get(memory.memory_id as usize)?.lock().await.take()?;
229 Some(T::from_memory_backend(backend).await)
230 }
231
232 pub async fn close_memory<T: FromMemoryBackend>(&self, device: T) -> Result<()> {
238 let backend = device.close_memory();
239 if let Some(mutex) = self.backends.get(backend.memory_id as usize) {
240 let mut guard = mutex.lock().await;
241 if guard.is_none() {
242 *guard = Some(backend);
243 } else {
244 println!("Warning: Attempted to close memory ID {} which is already closed", backend.memory_id);
245 }
246 } else {
247 println!("Warning: Attempted to close memory ID {} which does not exist", backend.memory_id);
248 }
249 Ok(())
250 }
251
252 pub async fn initialize_memory<T: FromMemoryBackend>(&self, memory: MemoryDevice) -> Option<Result<T>> {
260 let backend = self.backends.get(memory.memory_id as usize)?.lock().await.take()?;
261 Some(T::initialize_memory_backend(backend).await)
262 }
263
264}