From 6b27e84e4bee3a6318ed26e32871fc087d9e35b1 Mon Sep 17 00:00:00 2001 From: Jacob Weaver Date: Fri, 13 Dec 2024 10:23:45 -0600 Subject: [PATCH] firstadd --- .cargo/config.toml | 0 .gitignore | 21 ++ .vscode/settings.json | 3 + Cargo.toml | 24 ++ apps/ds/Cargo.toml | 19 ++ apps/ds/src/ds.rs | 192 +++++++++++ apps/ds/src/ds_std.rs | 72 ++++ apps/ds/src/ds_stub.rs | 25 ++ apps/example/Cargo.toml | 12 + apps/example/src/example.rs | 69 ++++ apps/hs/Cargo.toml | 22 ++ apps/hs/src/hs.rs | 154 +++++++++ apps/hs/src/hs_std.rs | 59 ++++ apps/hs/src/watchdog.rs | 84 +++++ apps/to/Cargo.toml | 13 + apps/to/src/to.rs | 151 +++++++++ builds/example_build/Cargo.toml | 15 + builds/example_build/src/main.rs | 140 ++++++++ builds/ground/Cargo.toml | 10 + builds/ground/src/main.rs | 30 ++ builds/rp_pico_example/.cargo/config.toml | 6 + builds/rp_pico_example/Cargo.toml | 39 +++ builds/rp_pico_example/build.rs | 25 ++ builds/rp_pico_example/build.sh | 3 + builds/rp_pico_example/memory.x | 15 + builds/rp_pico_example/src/main.rs | 204 ++++++++++++ rfe/Cargo.toml | 17 + rfe/macros/Cargo.toml | 14 + rfe/macros/src/macros.rs | 169 ++++++++++ rfe/src/connector.rs | 187 +++++++++++ rfe/src/lib.rs | 30 ++ rfe/src/msg.rs | 222 +++++++++++++ rfe/src/rfe.rs | 388 ++++++++++++++++++++++ rfe/src/time.rs | 45 +++ rfe/src/to_csv.rs | 151 +++++++++ rfe/src/utils.rs | 63 ++++ tools/decom/Cargo.toml | 15 + tools/decom/src/decom.rs | 113 +++++++ 38 files changed, 2821 insertions(+) create mode 100644 .cargo/config.toml create mode 100644 .gitignore create mode 100644 .vscode/settings.json create mode 100644 Cargo.toml create mode 100644 apps/ds/Cargo.toml create mode 100644 apps/ds/src/ds.rs create mode 100644 apps/ds/src/ds_std.rs create mode 100644 apps/ds/src/ds_stub.rs create mode 100644 apps/example/Cargo.toml create mode 100644 apps/example/src/example.rs create mode 100644 apps/hs/Cargo.toml create mode 100644 apps/hs/src/hs.rs create mode 100644 apps/hs/src/hs_std.rs create mode 100644 apps/hs/src/watchdog.rs create mode 100644 apps/to/Cargo.toml create mode 100644 apps/to/src/to.rs create mode 100644 builds/example_build/Cargo.toml create mode 100644 builds/example_build/src/main.rs create mode 100644 builds/ground/Cargo.toml create mode 100644 builds/ground/src/main.rs create mode 100644 builds/rp_pico_example/.cargo/config.toml create mode 100644 builds/rp_pico_example/Cargo.toml create mode 100644 builds/rp_pico_example/build.rs create mode 100755 builds/rp_pico_example/build.sh create mode 100644 builds/rp_pico_example/memory.x create mode 100644 builds/rp_pico_example/src/main.rs create mode 100644 rfe/Cargo.toml create mode 100644 rfe/macros/Cargo.toml create mode 100644 rfe/macros/src/macros.rs create mode 100644 rfe/src/connector.rs create mode 100644 rfe/src/lib.rs create mode 100644 rfe/src/msg.rs create mode 100644 rfe/src/rfe.rs create mode 100644 rfe/src/time.rs create mode 100644 rfe/src/to_csv.rs create mode 100644 rfe/src/utils.rs create mode 100644 tools/decom/Cargo.toml create mode 100644 tools/decom/src/decom.rs diff --git a/.cargo/config.toml b/.cargo/config.toml new file mode 100644 index 0000000..e69de29 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..64bf0eb --- /dev/null +++ b/.gitignore @@ -0,0 +1,21 @@ +# Generated by Cargo +# will have compiled files and executables +debug/ +target/ + +# Remove Cargo.lock from gitignore if creating an executable, leave it for libraries +# More information here https://doc.rust-lang.org/cargo/guide/cargo-toml-vs-cargo-lock.html +Cargo.lock + +# These are backup files generated by rustfmt +**/*.rs.bk + +# MSVC Windows builds of rustc generate these, which store debugging information +*.pdb + +# RustRover +# JetBrains specific template is maintained in a separate JetBrains.gitignore that can +# be found at https://github.com/github/gitignore/blob/main/Global/JetBrains.gitignore +# and can be added to the global gitignore or merged into this file. For a more nuclear +# option (not recommended) you can uncomment the following to ignore the entire idea folder. +#.idea/ diff --git a/.vscode/settings.json b/.vscode/settings.json new file mode 100644 index 0000000..23fd35f --- /dev/null +++ b/.vscode/settings.json @@ -0,0 +1,3 @@ +{ + "editor.formatOnSave": true +} \ No newline at end of file diff --git a/Cargo.toml b/Cargo.toml new file mode 100644 index 0000000..96068ef --- /dev/null +++ b/Cargo.toml @@ -0,0 +1,24 @@ +[workspace] + +members = [ + "apps/ds", + "apps/example", + "apps/hs", + "apps/to", + "builds/example_build", + "builds/ground", + "builds/rp_pico_example", + "rfe", + "tools/decom", +] +resolver = "2" + +[workspace.dependencies] +anyhow = { version = "1.0.93", default-features = false } +log = { version = "0.4.22", default-features = false } +hashbrown = { version = "0.15.1", features = ["alloc"] } +bincode = { version = "2.0.0-rc.3", default-features = false, features = [ + "alloc", + "derive", +] } +simple_logger = "5.0.0" diff --git a/apps/ds/Cargo.toml b/apps/ds/Cargo.toml new file mode 100644 index 0000000..e937e1f --- /dev/null +++ b/apps/ds/Cargo.toml @@ -0,0 +1,19 @@ +[package] +name = "ds" +version = "0.1.0" +edition = "2021" + +[lib] +path = "src/ds.rs" + +[features] +default = ["std"] +std = [] + +[dependencies] +rfe = { path = "../../rfe" } +anyhow.workspace = true +log.workspace = true +chrono = "0.4.38" +bincode.workspace = true +hashbrown.workspace = true diff --git a/apps/ds/src/ds.rs b/apps/ds/src/ds.rs new file mode 100644 index 0000000..0037005 --- /dev/null +++ b/apps/ds/src/ds.rs @@ -0,0 +1,192 @@ +#![no_std] + +#[cfg(feature = "std")] +mod ds_std; +use bincode::encode_to_vec; +#[cfg(feature = "std")] +pub use ds_std::DsFile; + +#[cfg(not(feature = "std"))] +mod ds_stub; +#[cfg(not(feature = "std"))] +pub use ds_stub::DsFile; + +extern crate alloc; + +use anyhow::Result; +use hashbrown::HashMap; +use log::*; +use msg::{DsCmd, DsHk, DsOutData, DsTlmSet, Instance, Msg, TlmSetId}; +use rfe::*; + +#[derive(Debug, Default)] +pub struct DsData { + hk: DsHk, + out_data: DsOutData, + file_list: HashMap, + enabled: bool, +} + +pub struct DsFileSettings { + pub max_size: u32, + pub max_age: u32, + pub enabled: bool, +} + +pub struct Ds { + data: DsData, + tlm_sets: HashMap, + start_enabled: bool, +} + +impl Ds { + pub fn new(tlm_sets: HashMap, start_enabled: bool) -> Self { + Self { + data: Default::default(), + tlm_sets, + start_enabled, + } + } + + pub fn update_subscriptions(&mut self, rfe: &mut Rfe) { + rfe.unsubscribe_all(); + rfe.subscribe_all( + self.tlm_sets + .values() + .filter(|x| x.enabled) + .map(|x| &x.items) + .flatten() + .map(|x| x.target), + ); + } +} + +impl App for Ds { + fn init(&mut self, rfe: &mut rfe::Rfe) -> Result<()> { + self.data = Default::default(); + self.data.enabled = self.start_enabled; + + self.update_subscriptions(rfe); + return Ok(()); + } + + fn run(&mut self, rfe: &mut rfe::Rfe) { + self.data.out_data.counter += 1; + self.data.out_data.bytes_written_this_cycle = 0; + while let Some(msg) = rfe.recv() { + match msg.msg { + Msg::DsCmd(cmd) => match cmd { + DsCmd::Noop => info!("Noop command received"), + DsCmd::Reset => { + info!("Reset command received"); + self.data = Default::default(); + } + DsCmd::CloseAll => { + info!("CloseAll command received"); + for f in self.data.file_list.values_mut() { + f.close(); + } + } + DsCmd::Close(f) => { + info!("Close command received"); + if let Some(file) = self.data.file_list.get_mut(&f) { + file.close(); + } else { + error!("Cannot close file {f}, file doesn't exist"); + } + } + DsCmd::AddTlmSet(ds_tlm_set) => { + info!("received AddTlmSet"); + if let Err(e) = self.tlm_sets.try_insert(ds_tlm_set.id, ds_tlm_set.clone()) + { + error!("Could not add tlm set {} {e}", ds_tlm_set.id); + } else { + info!("TlmSet {} added", ds_tlm_set.id); + self.update_subscriptions(rfe); + } + } + DsCmd::RemoveTlmSet(set_id) => { + info!("received RemoveTlmSet"); + if let Some(_set) = self.tlm_sets.remove(&set_id) { + info!("removed tlm set {}", set_id); + self.update_subscriptions(rfe); + } else { + warn!("Cannot remove tlm set {}, does not exist", set_id); + } + } + DsCmd::DisableTlmSet(set_id) => { + info!("received DisableTlmSet"); + if let Some(set) = self.tlm_sets.get_mut(&set_id) { + info!("set {set_id} is now disabled"); + set.enabled = false; + self.update_subscriptions(rfe); + } else { + warn!("could not disable set {set_id}, does not exist"); + } + } + DsCmd::EnablTlmSet(set_id) => { + info!("received EnablTlmSet"); + if let Some(set) = self.tlm_sets.get_mut(&set_id) { + info!("set {set_id} is now enabled"); + set.enabled = true; + self.update_subscriptions(rfe); + } else { + warn!("could not enable set {set_id}, does not exist"); + } + } + }, + _ => { + if !self.data.enabled { + continue; + } + + for tlm_set in self.tlm_sets.values_mut().filter(|x| x.enabled) { + for item in &mut tlm_set.items { + let msg_target = msg.to_target(); + if msg_target == item.target + || (msg_target.msg == item.target.msg + && item.target.instance == Instance::All) + { + if item.counter % (item.decimation + 1) == 0 { + let file = + if let Some(f) = self.data.file_list.get_mut(&tlm_set.id) { + f + } else { + let f = DsFile::new(tlm_set.path.clone()); + self.data.file_list.insert(tlm_set.id, f); + self.data.file_list.get_mut(&tlm_set.id).unwrap() + }; + + let bytes = encode_to_vec(&msg, BINCODE_CONFIG) + .expect("failed serialize ds packet"); + if let Err(e) = file.write(&bytes) { + error!("file write error: {e}"); + } else { + self.data.out_data.bytes_written_this_cycle += + bytes.len() as u32; + } + } + item.counter += 1; + } + } + } + } + } + } + + self.data.out_data.bytes_written += self.data.out_data.bytes_written_this_cycle; + } + + fn hk(&mut self, rfe: &mut rfe::Rfe) { + self.data.hk.counter = self.data.out_data.counter; + rfe.send(Msg::DsHk(self.data.hk)); + } + + fn out_data(&mut self, rfe: &mut rfe::Rfe) { + rfe.send(Msg::DsOutData(self.data.out_data)); + } + + fn get_app_rate(&self) -> Rate { + Rate::Hz1 + } +} diff --git a/apps/ds/src/ds_std.rs b/apps/ds/src/ds_std.rs new file mode 100644 index 0000000..a6b1069 --- /dev/null +++ b/apps/ds/src/ds_std.rs @@ -0,0 +1,72 @@ +extern crate std; + +use alloc::{ + format, + string::{String, ToString}, +}; +use chrono::Utc; +use log::*; +use std::{fs::File, io::Write, path::Path}; + +#[derive(Debug)] +pub struct DsFile { + pub dir: String, + pub file: Option, + prefix: String, +} + +impl DsFile { + pub fn new(dir: String) -> Self { + Self { + file: None, + prefix: dir + .clone() + .split("/") + .last() + .unwrap_or("unnamed") + .to_string(), + dir, + } + } + + pub fn close(&mut self) { + self.file = None; + } + + pub fn open(&mut self) { + let time = Self::get_time(&self.prefix); + let file_path = Path::new(&self.dir).join(time); + self.file = match File::create(&file_path) { + Ok(f) => Some(f), + Err(e) => { + error!("failed to create file at {:?} {}", file_path, e); + None + } + }; + } + + fn get_time(prefix: &str) -> String { + let date = Utc::now(); + format!("{}_{}", prefix, date.format("%Y-%m-%d_%H-%M-%S.dat")) + } + + pub fn write(&mut self, buf: &[u8]) -> std::io::Result { + if self.file.is_none() { + self.open(); + } + + if let Some(f) = &mut self.file { + return f.write(buf); + } else { + return Ok(0); + } + } + + pub fn flush(&mut self) -> std::io::Result<()> { + if let Some(f) = &mut self.file { + return f.flush(); + } else { + return Ok(()); + } + } +} diff --git a/apps/ds/src/ds_stub.rs b/apps/ds/src/ds_stub.rs new file mode 100644 index 0000000..ef2a6ec --- /dev/null +++ b/apps/ds/src/ds_stub.rs @@ -0,0 +1,25 @@ +use alloc::string::String; +use anyhow::Result; + +#[derive(Debug)] +pub struct DsFile { + pub dir: String, +} + +impl DsFile { + pub fn new(dir: String) -> Self { + Self { dir } + } + + pub fn close(&mut self) {} + + pub fn open(&mut self) {} + + pub fn write(&mut self, _buf: &[u8]) -> Result { + return Ok(0); + } + + pub fn flush(&mut self) -> Result<()> { + return Ok(()); + } +} diff --git a/apps/example/Cargo.toml b/apps/example/Cargo.toml new file mode 100644 index 0000000..50aaedc --- /dev/null +++ b/apps/example/Cargo.toml @@ -0,0 +1,12 @@ +[package] +name = "example" +version = "0.1.0" +edition = "2021" + +[lib] +path = "src/example.rs" + +[dependencies] +rfe = { path = "../../rfe" } +anyhow.workspace = true +log.workspace = true diff --git a/apps/example/src/example.rs b/apps/example/src/example.rs new file mode 100644 index 0000000..4a10507 --- /dev/null +++ b/apps/example/src/example.rs @@ -0,0 +1,69 @@ +#![no_std] + +use anyhow::Result; +use log::info; +use msg::{ExampleCmd, ExampleHk, ExampleOutData, Instance, Msg, MsgKind, TargetMsg}; +use rfe::*; + +#[derive(Debug, Clone, Copy, Default)] +pub struct ExampleData { + hk: ExampleHk, + out_data: ExampleOutData, +} + +pub struct Example { + data: ExampleData, +} + +impl Example { + pub fn new() -> Self { + Self { + data: Default::default(), + } + } +} + +impl App for Example { + fn init(&mut self, rfe: &mut rfe::Rfe) -> Result<()> { + self.data = Default::default(); + rfe.subscribe(TargetMsg::new(Instance::Other, MsgKind::ExampleHk)); + rfe.subscribe(TargetMsg::new(rfe.get_instance(), MsgKind::ExampleCmd)); + return Ok(()); + } + + fn run(&mut self, rfe: &mut rfe::Rfe) { + self.data.out_data.counter += 1; + info!("example running {:?}", self.data); + while let Some(msg) = rfe.recv() { + match msg.msg { + Msg::ExampleCmd(cmd) => match cmd { + ExampleCmd::Noop => info!("NOOP command received"), + ExampleCmd::Reset => { + info!("RESET command received"); + self.data = Default::default(); + } + }, + _ => { + info!("example got msg {:?}", msg); + } + } + } + + if self.data.out_data.counter > 10 { + rfe.send(Msg::ExampleCmd(ExampleCmd::Reset)); + }; + } + + fn hk(&mut self, rfe: &mut rfe::Rfe) { + self.data.hk.counter = self.data.out_data.counter; + rfe.send(Msg::ExampleHk(self.data.hk)); + } + + fn out_data(&mut self, rfe: &mut rfe::Rfe) { + rfe.send(Msg::ExampleOutData(self.data.out_data)); + } + + fn get_app_rate(&self) -> Rate { + Rate::Hz1 + } +} diff --git a/apps/hs/Cargo.toml b/apps/hs/Cargo.toml new file mode 100644 index 0000000..54a6645 --- /dev/null +++ b/apps/hs/Cargo.toml @@ -0,0 +1,22 @@ +[package] +name = "hs" +version = "0.1.0" +edition = "2021" + +[lib] +path = "src/hs.rs" + +[features] +default = ["std"] +std = ["dep:sysinfo", "dep:watchdog-device"] + +[dependencies] +rfe = { path = "../../rfe" } +anyhow.workspace = true +log.workspace = true +watchdog-device = { version = "0.2.0", optional = true } +sysinfo = { version = "0.32.0", default-features = false, optional = true, features = [ + "component", + "disk", + "system", +] } diff --git a/apps/hs/src/hs.rs b/apps/hs/src/hs.rs new file mode 100644 index 0000000..3ac02ed --- /dev/null +++ b/apps/hs/src/hs.rs @@ -0,0 +1,154 @@ +#![no_std] +use anyhow::Result; +use log::*; +use msg::{HsCmd, HsHk, HsOutData, Msg, MsgKind, TargetMsg}; +use rfe::*; + +extern crate alloc; +use alloc::vec::Vec; + +mod watchdog; +use utils::ManualAuto; +pub use watchdog::*; + +#[cfg(feature = "std")] +mod hs_std; +#[cfg(feature = "std")] +pub use hs_std::*; + +#[cfg(not(any(feature = "std")))] +mod hs_stub; +#[cfg(not(any(feature = "std")))] +pub use hs_stub::*; + +pub trait SystemInfoGrabber { + fn check_cpu_usage(&mut self) -> Vec; + fn check_mem_usage(&mut self) -> u8; + fn check_fs_usage(&mut self) -> Vec; + fn check_temps(&mut self) -> Vec; +} + +#[derive(Debug, Default, Clone)] +pub struct HsData { + out_data: HsOutData, + hk: HsHk, +} + +#[derive(Debug, Clone, Copy)] +pub struct HsConfig { + pub cpu_checks: bool, + pub mem_checks: bool, + pub fs_checks: bool, + pub temp_checks: bool, + pub watchdog_enable: bool, + pub watchdog_timeout: i32, +} + +pub struct Hs<'a> { + data: HsData, + config: HsConfig, + grabber: &'a mut dyn SystemInfoGrabber, + wd: WatchdogRef<'a>, + wd_value: ManualAuto, +} + +impl<'a> Hs<'a> { + pub fn new( + config: HsConfig, + grabber: &'a mut dyn SystemInfoGrabber, + mut watchdog: WatchdogRef<'a>, + ) -> Self { + if config.watchdog_enable { + watchdog.enable(); + } else { + watchdog.disable(); + } + watchdog.set_time(config.watchdog_timeout); + Self { + data: Default::default(), + config, + grabber, + wd: watchdog, + wd_value: ManualAuto::new(config.watchdog_enable, false), + } + } + + fn reset(&mut self) { + self.data = Default::default(); + } +} + +impl App for Hs<'_> { + fn init(&mut self, rfe: &mut Rfe) -> Result<()> { + self.reset(); + rfe.subscribe(TargetMsg::new(rfe.get_instance(), MsgKind::HsCmd)); + + return Ok(()); + } + + fn run(&mut self, rfe: &mut Rfe) { + self.data.out_data.counter += 1; + while let Some(msg) = rfe.recv() { + match msg.msg { + Msg::HsCmd(cmd) => { + self.data.hk.cmd_counter += 1; + match cmd { + HsCmd::Noop => { + info!("Noop command received"); + } + HsCmd::Reset => { + info!("Reset command received"); + self.data = Default::default(); + } + HsCmd::WatchdogEnableManual(v) => self.wd_value.manual_set(v), + HsCmd::WatchdogEnableAuto(v) => self.wd_value.auto_set(v), + HsCmd::WatchdogResumeAuto => self.wd_value.resume_auto(), + } + } + _ => { + warn!( + "HS received unexpected message: {:?} from {:?}", + msg.msg.kind(), + msg.instance + ); + } + } + } + + if self.wd_value.has_changed() { + if *self.wd_value.get() { + self.wd.enable(); + } else { + self.wd.disable(); + } + } + self.wd.feed(); + + if self.config.cpu_checks { + self.data.hk.cpu_usage = self.grabber.check_cpu_usage(); + } + if self.config.mem_checks { + self.data.hk.mem_usage = self.grabber.check_mem_usage(); + } + if self.config.fs_checks { + self.data.hk.fs_usage = self.grabber.check_fs_usage(); + } + if self.config.temp_checks { + self.data.hk.temps = self.grabber.check_temps(); + } + } + + fn hk(&mut self, rfe: &mut Rfe) { + self.data.hk.counter = self.data.out_data.counter; + + rfe.send(Msg::HsHk(self.data.hk.clone())); + } + + fn out_data(&mut self, rfe: &mut Rfe) { + rfe.send(Msg::HsOutData(self.data.out_data)); + } + + fn get_app_rate(&self) -> Rate { + Rate::Hz1 + } +} diff --git a/apps/hs/src/hs_std.rs b/apps/hs/src/hs_std.rs new file mode 100644 index 0000000..8bc8196 --- /dev/null +++ b/apps/hs/src/hs_std.rs @@ -0,0 +1,59 @@ +use sysinfo::{Components, CpuRefreshKind, Disks, MemoryRefreshKind, RefreshKind, System}; + +use crate::SystemInfoGrabber; +extern crate alloc; +use alloc::vec::Vec; + +pub struct StdSystemInfoGrabber { + system: System, +} + +impl StdSystemInfoGrabber { + pub fn new() -> Self { + Self { + system: System::new_with_specifics( + RefreshKind::new() + .with_cpu(CpuRefreshKind::new().with_cpu_usage()) + .with_memory(MemoryRefreshKind::new().with_ram()), + ), + } + } +} + +impl SystemInfoGrabber for StdSystemInfoGrabber { + fn check_cpu_usage(&mut self) -> Vec { + self.system.refresh_cpu_usage(); + self.system + .cpus() + .iter() + .map(|x| x.cpu_usage().round() as u8) + .collect() + } + + fn check_mem_usage(&mut self) -> u8 { + self.system.refresh_memory(); + (self.system.used_memory() * 100 / self.system.total_memory()) as u8 + } + + fn check_fs_usage(&mut self) -> Vec { + let mut disks = Disks::new_with_refreshed_list(); + let disks = disks.list_mut(); + disks.sort_by(|x, y| x.name().cmp(y.name())); + + disks + .iter() + .map(|x| ((x.total_space() - x.available_space()) * 100 / x.total_space()) as u8) + .collect() + } + + fn check_temps(&mut self) -> Vec { + let mut components = Components::new_with_refreshed_list(); + let components = components.list_mut(); + components.sort_by(|x, y| x.label().cmp(y.label())); + + components + .iter() + .map(|x| x.temperature().round() as i8) + .collect() + } +} diff --git a/apps/hs/src/watchdog.rs b/apps/hs/src/watchdog.rs new file mode 100644 index 0000000..b50ae76 --- /dev/null +++ b/apps/hs/src/watchdog.rs @@ -0,0 +1,84 @@ +pub trait Watchdog { + fn enable(&mut self); + fn disable(&mut self); + fn set_time(&mut self, time: i32); + fn feed(&mut self); +} + +#[cfg(feature = "std")] +mod watchdog_std { + use super::Watchdog; + extern crate std; + use crate::unwrap_print_err; + use anyhow::Result; + use log::*; + + pub struct LinuxWatchdog { + wd: watchdog_device::Watchdog, + } + + impl LinuxWatchdog { + pub fn new() -> Result { + let wd = watchdog_device::Watchdog::new()?; + wd.set_option(&watchdog_device::SetOptionFlags::DisableCard)?; + Ok(Self { wd }) + } + } + + impl Watchdog for LinuxWatchdog { + fn enable(&mut self) { + unwrap_print_err!( + self.wd + .set_option(&watchdog_device::SetOptionFlags::EnableCard), + "failed to enable watchdog" + ); + } + + fn disable(&mut self) { + unwrap_print_err!( + self.wd + .set_option(&watchdog_device::SetOptionFlags::DisableCard), + "failed to set watchdog timeout" + ); + } + + fn set_time(&mut self, time: i32) { + unwrap_print_err!(self.wd.set_timeout(time), "failed to set watchdog timeout"); + } + + fn feed(&mut self) { + unwrap_print_err!(self.wd.keep_alive(), "failed to feed watchdog"); + } + } +} + +#[cfg(feature = "std")] +pub use watchdog_std::*; + +pub type WatchdogRef<'a> = Option<&'a mut dyn Watchdog>; + +impl<'a> Watchdog for Option<&'a mut dyn Watchdog> { + fn enable(&mut self) { + if let Some(wd) = self { + wd.enable(); + } + } + + fn disable(&mut self) { + if let Some(wd) = self { + wd.disable(); + } + } + + fn set_time(&mut self, time: i32) { + if let Some(wd) = self { + wd.set_time(time); + } + } + + fn feed(&mut self) { + if let Some(wd) = self { + wd.feed(); + } + } +} diff --git a/apps/to/Cargo.toml b/apps/to/Cargo.toml new file mode 100644 index 0000000..be09c9a --- /dev/null +++ b/apps/to/Cargo.toml @@ -0,0 +1,13 @@ +[package] +name = "to" +version = "0.1.0" +edition = "2021" + +[lib] +path = "src/to.rs" + +[dependencies] +rfe = { path = "../../rfe" } +anyhow.workspace = true +log.workspace = true +hashbrown.workspace = true diff --git a/apps/to/src/to.rs b/apps/to/src/to.rs new file mode 100644 index 0000000..313e63e --- /dev/null +++ b/apps/to/src/to.rs @@ -0,0 +1,151 @@ +#![no_std] +extern crate alloc; +use alloc::vec::Vec; +use anyhow::Result; +use connector::Connector; +use hashbrown::HashMap; +use log::*; +use msg::{Instance, Msg, MsgKind, TargetMsg, TlmSetId, ToCmd, ToHk, ToOutData, ToTlmSet}; +use rfe::*; + +#[derive(Debug, Clone, Default)] +pub struct ToData { + out_data: ToOutData, + hk: ToHk, +} + +pub struct To<'a> { + data: ToData, + connector: &'a mut dyn Connector, + tlm_sets: HashMap, +} + +impl<'a> To<'a> { + pub fn new(connector: &'a mut dyn Connector, tlm_sets: HashMap) -> Self { + Self { + connector, + data: Default::default(), + tlm_sets, + } + } + + pub fn update_subscriptions(&mut self, rfe: &mut Rfe) { + rfe.unsubscribe_all(); + rfe.subscribe_all( + self.tlm_sets + .values() + .filter(|x| x.enabled) + .map(|x| &x.items) + .flatten() + .map(|x| x.target), + ); + } + + pub fn handle_cmd(&mut self, rfe: &mut Rfe, cmd: &ToCmd) { + match cmd { + ToCmd::Noop => info!("received Noop"), + ToCmd::Reset => { + info!("received Reset"); + self.data = Default::default(); + } + ToCmd::AddTlmSet(to_tlm_set) => { + info!("received AddTlmSet"); + if let Err(e) = self.tlm_sets.try_insert(to_tlm_set.id, to_tlm_set.clone()) { + error!("Could not add tlm set {} {e}", to_tlm_set.id); + } else { + info!("TlmSet {} added", to_tlm_set.id); + self.update_subscriptions(rfe); + } + } + ToCmd::RemoveTlmSet(set_id) => { + info!("received RemoveTlmSet"); + if let Some(_set) = self.tlm_sets.remove(set_id) { + info!("removed tlm set {}", set_id); + self.update_subscriptions(rfe); + } else { + warn!("Cannot remove tlm set {}, does not exist", set_id); + } + } + ToCmd::DisableTlmSet(set_id) => { + info!("received DisableTlmSet"); + if let Some(set) = self.tlm_sets.get_mut(set_id) { + info!("set {set_id} is now disabled"); + set.enabled = false; + self.update_subscriptions(rfe); + } else { + warn!("could not disable set {set_id}, does not exist"); + } + } + ToCmd::EnablTlmSet(set_id) => { + info!("received EnablTlmSet"); + if let Some(set) = self.tlm_sets.get_mut(set_id) { + info!("set {set_id} is now enabled"); + set.enabled = true; + self.update_subscriptions(rfe); + } else { + warn!("could not enable set {set_id}, does not exist"); + } + } + } + info!("got cmd {:?}", cmd); + } +} + +impl App for To<'_> { + fn init(&mut self, rfe: &mut Rfe) -> Result<()> { + rfe.subscribe(TargetMsg::new(rfe.get_instance(), MsgKind::ToCmd)); + self.update_subscriptions(rfe); + return Ok(()); + } + + fn run(&mut self, rfe: &mut Rfe) { + self.data.out_data.counter += 1; + let mut msgs = Vec::new(); + while let Some(msg) = rfe.recv() { + let mut is_cmd = false; + if let Msg::ToCmd(cmd) = &msg.msg { + if msg.instance == rfe.get_instance() { + is_cmd = true; + self.handle_cmd(rfe, cmd); + } + } + if !is_cmd { + for tlm_set in self.tlm_sets.values_mut().filter(|x| x.enabled) { + for item in &mut tlm_set.items { + let msg_target = msg.to_target(); + if msg_target == item.target + || (msg_target.msg == item.target.msg + && item.target.instance == Instance::All) + { + if item.counter % (item.decimation + 1) == 0 { + msgs.push(msg.clone()); + } + item.counter += 1; + } + } + } + } + } + if msgs.len() > 0 { + self.connector.send(msgs); + } + + while let Some(msgs) = self.connector.recv() { + for msg in msgs { + rfe.post_message(msg); + } + } + } + + fn hk(&mut self, rfe: &mut Rfe) { + rfe.send(Msg::ToHk(self.data.hk)); + } + + fn out_data(&mut self, rfe: &mut Rfe) { + rfe.send(Msg::ToOutData(self.data.out_data)); + } + + fn get_app_rate(&self) -> Rate { + Rate::Hz50 + } +} diff --git a/builds/example_build/Cargo.toml b/builds/example_build/Cargo.toml new file mode 100644 index 0000000..1dba813 --- /dev/null +++ b/builds/example_build/Cargo.toml @@ -0,0 +1,15 @@ +[package] +name = "example_build" +version = "0.1.0" +edition = "2021" + +[dependencies] +simple_logger.workspace = true +rfe.path = "../../rfe" +rfe.features = ["std"] +example.path = "../../apps/example" +ds.path = "../../apps/ds" +hs.path = "../../apps/hs" +to.path = "../../apps/to" +anyhow.workspace = true +hashbrown.workspace = true diff --git a/builds/example_build/src/main.rs b/builds/example_build/src/main.rs new file mode 100644 index 0000000..a5a02d9 --- /dev/null +++ b/builds/example_build/src/main.rs @@ -0,0 +1,140 @@ +use anyhow::Result; +use connector::UdpConnector; +use ds::*; +use example::*; +use hashbrown::HashMap; +use hs::*; +use msg::{DsTlmSet, Instance, MsgKind, TargetMsg, TlmSetItem, ToTlmSet}; +use rfe::*; +use simple_logger::SimpleLogger; +use time::UnixTimeDriver; +use to::*; + +fn main() -> Result<()> { + SimpleLogger::new().init().unwrap(); + let mut record = HashMap::new(); + record.insert( + 0, + DsTlmSet { + enabled: true, + path: "log/example".to_string(), + items: vec![ + TlmSetItem { + counter: 0, + target: TargetMsg::new(Instance::All, MsgKind::ExampleOutData), + decimation: 0, + }, + TlmSetItem { + counter: 0, + target: TargetMsg::new(Instance::All, MsgKind::ExampleHk), + decimation: 0, + }, + ], + id: 0, + }, + ); + record.insert( + 1, + DsTlmSet { + enabled: true, + path: "log/ds".to_string(), + items: vec![ + TlmSetItem { + counter: 0, + target: TargetMsg::new(Instance::All, MsgKind::DsOutData), + decimation: 0, + }, + TlmSetItem { + counter: 0, + target: TargetMsg::new(Instance::All, MsgKind::DsHk), + decimation: 0, + }, + ], + id: 1, + }, + ); + record.insert( + 2, + DsTlmSet { + enabled: true, + path: "log/hs".to_string(), + items: vec![ + TlmSetItem { + counter: 0, + target: TargetMsg::new(Instance::All, MsgKind::HsOutData), + decimation: 0, + }, + TlmSetItem { + counter: 0, + target: TargetMsg::new(Instance::All, MsgKind::HsHk), + decimation: 0, + }, + ], + id: 2, + }, + ); + record.insert( + 3, + DsTlmSet { + enabled: true, + path: "log/to".to_string(), + items: vec![ + TlmSetItem { + counter: 0, + target: TargetMsg::new(Instance::All, MsgKind::ToOutData), + decimation: 0, + }, + TlmSetItem { + counter: 0, + target: TargetMsg::new(Instance::All, MsgKind::ToHk), + decimation: 0, + }, + ], + id: 3, + }, + ); + + let mut example = Example::new(); + let mut ds = Ds::new(record, false); + // let mut wd = LinuxWatchdog::new().unwrap(); + let mut grabber = StdSystemInfoGrabber::new(); + let mut hs = Hs::new( + HsConfig { + cpu_checks: true, + mem_checks: true, + fs_checks: true, + temp_checks: true, + watchdog_enable: true, + watchdog_timeout: 10, + }, + &mut grabber, + // Some(&mut wd), + None, + ); + let mut tcp = UdpConnector::new("127.0.0.1", 7412, "127.0.0.1", 7413)?; + let mut ground_connector = UdpConnector::new("127.0.0.1", 7010, "127.0.0.1", 7011)?; + let mut dl_sets = HashMap::new(); + dl_sets.insert( + 0, + ToTlmSet { + items: vec![TlmSetItem { + counter: 0, + decimation: 9, + target: TargetMsg::new(Instance::All, MsgKind::HsHk), + }], + id: 0, + enabled: true, + }, + ); + let mut to = To::new(&mut ground_connector, dl_sets); + let time_driver = UnixTimeDriver::new(); + let mut instance = RfeInstance::new(Instance::Example, &time_driver); + instance.add_app("example", &mut example)?; + instance.add_app("to", &mut to)?; + instance.add_app("DS", &mut ds)?; + instance.add_app("HS", &mut hs)?; + instance.add_connector(&mut tcp); + + instance.start(); + return Ok(()); +} diff --git a/builds/ground/Cargo.toml b/builds/ground/Cargo.toml new file mode 100644 index 0000000..6b5b03f --- /dev/null +++ b/builds/ground/Cargo.toml @@ -0,0 +1,10 @@ +[package] +name = "ground" +version = "0.1.0" +edition = "2021" + +[dependencies] +rfe.path = "../../rfe" +anyhow.workspace = true +simple_logger = "5.0.0" +log.workspace = true diff --git a/builds/ground/src/main.rs b/builds/ground/src/main.rs new file mode 100644 index 0000000..c0b6303 --- /dev/null +++ b/builds/ground/src/main.rs @@ -0,0 +1,30 @@ +#![feature(thread_sleep_until)] + +use std::{ + thread::sleep_until, + time::{Duration, Instant}, +}; + +use anyhow::Result; +use connector::{Connector, UdpConnector}; +use log::*; +use rfe::*; +use simple_logger::SimpleLogger; + +fn main() -> Result<()> { + SimpleLogger::new().init().unwrap(); + + let mut udp = UdpConnector::new("127.0.0.1", 7011, "127.0.0.1", 7010)?; + + let mut next_time = Instant::now() + Duration::from_millis(10); + loop { + sleep_until(next_time); + + while let Some(msgs) = udp.recv() { + for msg in msgs { + info!("got msg {:?}", msg); + } + } + next_time += Duration::from_millis(10); + } +} diff --git a/builds/rp_pico_example/.cargo/config.toml b/builds/rp_pico_example/.cargo/config.toml new file mode 100644 index 0000000..541616c --- /dev/null +++ b/builds/rp_pico_example/.cargo/config.toml @@ -0,0 +1,6 @@ +[target.'cfg(all(target_arch = "arm", target_os = "none"))'] +runner = "elf2uf2-rs -d" + + +[build] +target = "thumbv6m-none-eabi" diff --git a/builds/rp_pico_example/Cargo.toml b/builds/rp_pico_example/Cargo.toml new file mode 100644 index 0000000..8d27eb7 --- /dev/null +++ b/builds/rp_pico_example/Cargo.toml @@ -0,0 +1,39 @@ +[package] +name = "rp_pico_example" +version = "0.1.0" +edition = "2021" + +[dependencies] +embassy-usb-logger = { git = "https://github.com/embassy-rs/embassy.git" } +rfe = { path = "../../rfe", default-features = false } +to = { path = "../../apps/to", default-features = false } +anyhow.workspace = true +anyhow.default_features = false +embedded-alloc = "0.6.0" +embedded-hal = { version = "0.2.7", features = ["unproven"] } +embassy-rp = { git = "https://github.com/embassy-rs/embassy.git", default-features = false, features = [ + "rp2040", + "log", +] } +log.workspace = true +static_cell = "2.1.0" +portable-atomic = { version = "1.9.0", features = ["critical-section"] } +rtic = { version = "2.1.2", features = ["thumbv6-backend"] } +rtic-monotonics = { version = "2.0.3", features = ["rp2040"] } +rp-pico = { version = "0.9.0", default-features = false, features = [ + "rt", + "critical-section-impl", + "disable-intrinsics", +] } +panic-halt = "1.0.0" +hashbrown.workspace = true + +[profile.release] +debug = 2 +lto = true +opt-level = 'z' + +[profile.dev] +debug = 2 +lto = true +opt-level = "z" diff --git a/builds/rp_pico_example/build.rs b/builds/rp_pico_example/build.rs new file mode 100644 index 0000000..38c07f7 --- /dev/null +++ b/builds/rp_pico_example/build.rs @@ -0,0 +1,25 @@ +use std::env; +use std::fs::File; +use std::io::Write; +use std::path::PathBuf; + +fn main() { + // Put `memory.x` in our output directory and ensure it's + // on the linker search path. + let out = &PathBuf::from(env::var_os("OUT_DIR").unwrap()); + File::create(out.join("memory.x")) + .unwrap() + .write_all(include_bytes!("memory.x")) + .unwrap(); + println!("cargo:rustc-link-search={}", out.display()); + + // By default, Cargo will re-run a build script whenever + // any file in the project changes. By specifying `memory.x` + // here, we ensure the build script is only re-run when + // `memory.x` is changed. + println!("cargo:rerun-if-changed=memory.x"); + + println!("cargo:rustc-link-arg-bins=--nmagic"); + println!("cargo:rustc-link-arg-bins=-Tlink.x"); + // println!("cargo:rustc-link-arg-bins=-Tlink-rp.x"); +} diff --git a/builds/rp_pico_example/build.sh b/builds/rp_pico_example/build.sh new file mode 100755 index 0000000..cebe09c --- /dev/null +++ b/builds/rp_pico_example/build.sh @@ -0,0 +1,3 @@ +#!/bin/bash + +cargo run --release --target thumbv6m-none-eabi \ No newline at end of file diff --git a/builds/rp_pico_example/memory.x b/builds/rp_pico_example/memory.x new file mode 100644 index 0000000..070eac7 --- /dev/null +++ b/builds/rp_pico_example/memory.x @@ -0,0 +1,15 @@ +MEMORY { + BOOT2 : ORIGIN = 0x10000000, LENGTH = 0x100 + FLASH : ORIGIN = 0x10000100, LENGTH = 2048K - 0x100 + RAM : ORIGIN = 0x20000000, LENGTH = 256K +} + +EXTERN(BOOT2_FIRMWARE) + +SECTIONS { + /* ### Boot loader */ + .boot2 ORIGIN(BOOT2) : + { + KEEP(*(.boot2)); + } > BOOT2 +} INSERT BEFORE .text; \ No newline at end of file diff --git a/builds/rp_pico_example/src/main.rs b/builds/rp_pico_example/src/main.rs new file mode 100644 index 0000000..e02e886 --- /dev/null +++ b/builds/rp_pico_example/src/main.rs @@ -0,0 +1,204 @@ +#![no_std] +#![no_main] + +use rtic_monotonics::rp2040::prelude::*; + +rp2040_timer_monotonic!(Mono); + +#[rtic::app(device = rp_pico::hal::pac, dispatchers = [SW0_IRQ])] +mod app { + use super::*; + + use embassy_rp::{ + interrupt::typelevel, + peripherals::USB, + usb::{self, Driver}, + }; + use embedded_hal::digital::v2::OutputPin; + use fugit::Duration; + use log::info; + use msg::{ExampleHk, Instance, MsgKind, MsgPacket, TargetMsg, TlmSetItem, ToTlmSet}; + use rfe::{connector::Connector, Rate, *}; + use rp_pico::hal::{ + clocks, + gpio::{bank0::Gpio25, FunctionSio, Pin, PullNone, SioOutput}, + Sio, Watchdog, + }; + + use anyhow::Result; + use core::{mem::MaybeUninit, ptr::addr_of_mut}; + use embedded_alloc::LlffHeap as Heap; + use panic_halt as _; + use rp_pico::XOSC_CRYSTAL_FREQ; + use time::SchTimeDriver; + extern crate alloc; + use alloc::vec; + use alloc::vec::Vec; + use hashbrown::HashMap; + use to::*; + + #[global_allocator] + static HEAP: Heap = Heap::empty(); + + /// will log messages to log + #[derive(Debug)] + struct LogConnector; + + impl Connector for LogConnector { + fn send(&mut self, msgs: Vec) { + info!("got msgs: {:?}", msgs); + } + fn recv(&mut self) -> Option> { + return None; + } + } + + struct BlinkApp<'a> { + led_pin: &'a mut Pin, PullNone>, + on: bool, + counter: u32, + } + + impl App for BlinkApp<'_> { + fn init(&mut self, _rfe: &mut Rfe) -> Result<()> { + return Ok(()); + } + + fn run(&mut self, _rfe: &mut Rfe) { + self.counter += 1; + self.on = !self.on; + if self.on { + info!("led on!"); + self.led_pin.set_high().ok(); + } else { + info!("led off!"); + self.led_pin.set_low().ok(); + } + } + + fn hk(&mut self, rfe: &mut Rfe) { + rfe.send(msg::Msg::ExampleHk(ExampleHk { + counter: self.counter, + })); + } + + fn out_data(&mut self, _rfe: &mut Rfe) {} + + fn get_app_rate(&self) -> Rate { + return Rate::Hz10; + } + } + + #[shared] + struct Shared {} + + #[local] + struct Local { + led: Pin, PullNone>, + usb: Option>, + } + + struct UsbBinding; + unsafe impl typelevel::Binding> for UsbBinding {} + + #[init()] + fn init(mut ctx: init::Context) -> (Shared, Local) { + let p = embassy_rp::init(Default::default()); + // Configure the clocks, watchdog - The default is to generate a 125 MHz system clock + Mono::start(ctx.device.TIMER, &mut ctx.device.RESETS); // default rp2040 clock-rate is 125MHz + let mut watchdog = Watchdog::new(ctx.device.WATCHDOG); + let _clocks = clocks::init_clocks_and_plls( + XOSC_CRYSTAL_FREQ, + ctx.device.XOSC, + ctx.device.CLOCKS, + ctx.device.PLL_SYS, + ctx.device.PLL_USB, + &mut ctx.device.RESETS, + &mut watchdog, + ) + .ok() + .unwrap(); + + { + const HEAP_SIZE: usize = 1024 * 16; + static mut HEAP_MEM: [MaybeUninit; HEAP_SIZE] = [MaybeUninit::uninit(); HEAP_SIZE]; + unsafe { HEAP.init(addr_of_mut!(HEAP_MEM) as usize, HEAP_SIZE) } + } + + let driver = Driver::new(p.USB, UsbBinding {}); + let sio = Sio::new(ctx.device.SIO); + let gpioa = rp_pico::Pins::new( + ctx.device.IO_BANK0, + ctx.device.PADS_BANK0, + sio.gpio_bank0, + &mut ctx.device.RESETS, + ); + let led = gpioa + .led + .into_pull_type::() + .into_push_pull_output(); + + // Spawn heartbeat task + blink_instance::spawn().ok(); + usb_logger::spawn().ok(); + + // Return resources and timer + ( + Shared {}, + Local { + led, + usb: Some(driver), + }, + ) + } + + #[task(local = [usb])] + async fn usb_logger(ctx: usb_logger::Context) { + let driver = Option::take(ctx.local.usb).unwrap(); + embassy_usb_logger::run!(1024, log::LevelFilter::Info, driver); + } + + #[task(binds = USBCTRL_IRQ)] + fn usbctrl(_cx: usbctrl::Context) { + unsafe { + as typelevel::Handler>::on_interrupt(); + } + } + + #[task(local = [led], priority = 1)] + async fn blink_instance(mut ctx: blink_instance::Context) { + let mut blink_app = BlinkApp { + led_pin: &mut ctx.local.led, + on: false, + counter: 0, + }; + + let mut log_connector = LogConnector {}; + let mut tlmsets = HashMap::new(); + tlmsets.insert( + 0, + ToTlmSet { + items: vec![TlmSetItem { + target: TargetMsg::new(Instance::All, MsgKind::ExampleHk), + counter: 0, + decimation: 0, + }], + id: 0, + enabled: true, + }, + ); + let mut to = To::new(&mut log_connector, tlmsets); + let time_driver = SchTimeDriver::new(); + let mut instance = RfeInstance::new(Instance::Example, &time_driver); + instance.add_app("blink_app", &mut blink_app).unwrap(); + instance.add_app("to", &mut to).unwrap(); + + let mut next_time = Mono::now() + Duration::::from_ticks(10000); + + loop { + instance.run(); + Mono::delay_until(next_time).await; + next_time += Duration::::from_ticks(10000); + } + } +} diff --git a/rfe/Cargo.toml b/rfe/Cargo.toml new file mode 100644 index 0000000..6a787b9 --- /dev/null +++ b/rfe/Cargo.toml @@ -0,0 +1,17 @@ +[package] +name = "rfe" +version = "0.1.0" +edition = "2021" + +[features] +default = [] +std = ["dep:mio"] +to_csv = [] + +[dependencies] +bincode.workspace = true +anyhow.workspace = true +hashbrown.workspace = true +log.workspace = true +mio = { version = "1.0.2", features = ["net", "os-poll"], optional = true } +macros.path = "macros" diff --git a/rfe/macros/Cargo.toml b/rfe/macros/Cargo.toml new file mode 100644 index 0000000..9461a41 --- /dev/null +++ b/rfe/macros/Cargo.toml @@ -0,0 +1,14 @@ +[package] +name = "macros" +version = "0.1.0" +edition = "2021" + +[lib] +path = "src/macros.rs" +proc-macro = true + + +[dependencies] +proc-macro2 = "1.0.89" +quote = "1.0.37" +syn = { version = "2.0.87", features = ["full"] } diff --git a/rfe/macros/src/macros.rs b/rfe/macros/src/macros.rs new file mode 100644 index 0000000..8398bf5 --- /dev/null +++ b/rfe/macros/src/macros.rs @@ -0,0 +1,169 @@ +use proc_macro::TokenStream; +use quote::{format_ident, quote}; +use syn::Fields; +use syn::{parse_macro_input, Data, DeriveInput}; + +#[proc_macro_derive(Kind)] +pub fn kind_derive(input: TokenStream) -> TokenStream { + let input = parse_macro_input!(input as DeriveInput); + + let name = input.ident; + let kind_name = format_ident!("{name}Kind"); + + let variants = if let Data::Enum(data_enum) = input.data { + data_enum.variants + } else { + return syn::Error::new_spanned(name, "Kind derive can only be used on enums") + .to_compile_error() + .into(); + }; + + let kind_enum_arms = variants.iter().map(|variant| { + let variant_name = &variant.ident; + quote! { + #variant_name, + } + }); + let kind_enum_vec_arms = variants.iter().map(|variant| { + let variant_name = &variant.ident; + quote! { + #kind_name::#variant_name, + } + }); + + let kind_arms = variants.iter().map(|variant| { + let variant_name = &variant.ident; + match &variant.fields { + Fields::Unit => { + quote! { + Self::#variant_name => #kind_name::#variant_name, + } + } + Fields::Unnamed(_) => { + quote! { + Self::#variant_name(_) => #kind_name::#variant_name, + } + } + Fields::Named(_) => { + panic!("Named fields not supported by Kind"); + } + } + }); + + let expanded = quote! { + #[derive(Debug, Copy, Clone, PartialEq, Eq, Hash, Encode, Decode)] + #[cfg_attr(feature = "to_csv", derive(ToCsv))] + pub enum #kind_name { + #(#kind_enum_arms)* + } + + impl #kind_name { + pub fn to_vec() -> alloc::vec::Vec { + alloc::vec![ + #(#kind_enum_vec_arms)* + ] + } + } + + impl #name { + pub fn kind(&self) -> #kind_name { + match self { + #(#kind_arms)* + } + } + } + }; + + TokenStream::from(expanded) +} + +#[proc_macro_derive(ToCsv)] +pub fn to_csv_derive(input: TokenStream) -> TokenStream { + let input = parse_macro_input!(input as DeriveInput); + + let name = input.ident; + let name_s = name.to_string(); + + let expanded = match input.data { + Data::Struct(data_struct) => { + let fields = data_struct.fields; + + let arms = fields.iter().map(|field| { + let field_ident = field.clone().ident.unwrap(); + let field_name = field_ident.to_string(); + + quote! { { + let mut field_csvs: Vec = (&self.#field_ident as &dyn ToCsv).to_csv(); + for entry in &mut field_csvs { + *entry = format!("{}.{}", #field_name, &entry); + } + values.extend(field_csvs); + } } + }); + + quote! { + impl ToCsv for #name { + fn to_csv(&self) -> Vec { + extern crate alloc; + use alloc::format; + + let mut values:Vec = Vec::new(); + #(#arms)* + for value in &mut values { + *value = format!("{}", value); + } + return values; + } + } + } + } + Data::Enum(data_enum) => { + let variants = data_enum.variants; + let arms = variants.iter().map(|variant| { + let variant_name = &variant.ident; + let variant_name_s = variant_name.to_string(); + + match &variant.fields { + Fields::Unit => { + quote! { Self::#variant_name => values.push(format!(" = {}::{}", #name_s, #variant_name_s)), } + // quote! { Self::#variant_name => {}, } + } + Fields::Unnamed(_) => { + quote! {Self::#variant_name(l) => { + let mut field_csvs: Vec = (l as &dyn ToCsv).to_csv(); + for entry in &mut field_csvs { + *entry = format!("{}.{}", #variant_name_s, &entry); + } + values.extend(field_csvs); + },} + } + Fields::Named(_) => { + panic!("Named fields not supported by ToCsv"); + } + } + }); + + quote! { + impl ToCsv for #name { + fn to_csv(&self) -> Vec { + extern crate alloc; + use alloc::format; + let mut values: Vec = Vec::new(); + + match self { + #(#arms)* + } + return values; + } + } + } + } + Data::Union(_data_union) => { + return syn::Error::new_spanned(name, "ToCsv not implemented for unions") + .to_compile_error() + .into(); + } + }; + + TokenStream::from(expanded) +} diff --git a/rfe/src/connector.rs b/rfe/src/connector.rs new file mode 100644 index 0000000..a2e8a9b --- /dev/null +++ b/rfe/src/connector.rs @@ -0,0 +1,187 @@ +use core::fmt::Debug; + +use crate::msg::MsgPacket; +extern crate alloc; +use alloc::vec::Vec; + +pub trait Connector: Debug { + fn send(&mut self, msgs: Vec); + fn recv(&mut self) -> Option>; +} + +#[cfg(feature = "std")] +mod connector_std { + extern crate alloc; + extern crate std; + use alloc::vec::Vec; + use anyhow::{anyhow, Result}; + use bincode::{decode_from_slice, encode_to_vec}; + use core::time::Duration; + use log::*; + use mio::net::{TcpListener, TcpStream, UdpSocket}; + use std::{ + io::{Read, Write}, + net::ToSocketAddrs, + sync::mpsc::{self, Receiver, Sender}, + }; + + use super::Connector; + use crate::{msg::MsgPacket, BINCODE_CONFIG}; + + #[derive(Debug)] + pub struct MemConnector { + sender: Sender>, + receiver: Receiver>, + } + + impl MemConnector { + pub fn new() -> (Self, Self) { + let (s1, r1) = mpsc::channel(); + let (s2, r2) = mpsc::channel(); + return ( + Self { + sender: s1, + receiver: r2, + }, + Self { + sender: s2, + receiver: r1, + }, + ); + } + } + + impl Connector for MemConnector { + fn send(&mut self, msgs: Vec) { + self.sender.send(msgs).ok(); + } + + fn recv(&mut self) -> Option> { + self.receiver.recv_timeout(Duration::ZERO).ok() + } + } + + #[derive(Debug)] + pub struct TcpConnector { + listener: TcpListener, + listen: Option, + client: TcpStream, + remote_addr: &'static str, + remote_port: u16, + } + + impl TcpConnector { + pub fn new( + local_addr: &'static str, + local_port: u16, + remote_addr: &'static str, + remote_port: u16, + ) -> Result { + Ok(Self { + listener: TcpListener::bind( + (local_addr, local_port) + .to_socket_addrs()? + .next() + .ok_or(anyhow!("failed to parse ip address"))?, + )?, + client: TcpStream::connect( + (remote_addr, remote_port) + .to_socket_addrs()? + .next() + .ok_or(anyhow!("failed to parse ip address"))?, + )?, + listen: None, + remote_addr, + remote_port, + }) + } + } + + impl Connector for TcpConnector { + fn send(&mut self, msgs: Vec) { + let r = encode_to_vec(&msgs, BINCODE_CONFIG).expect("failed to serialize tcp packet"); + if let Err(e) = self.client.write(&r) { + warn!("tcp write error {e}"); + self.client = TcpStream::connect( + (self.remote_addr, self.remote_port) + .to_socket_addrs() + .unwrap() + .next() + .ok_or(anyhow!("failed to parse ip address")) + .unwrap(), + ) + .unwrap(); + } + } + + fn recv(&mut self) -> Option> { + let mut read_buf = [0_u8; 4096]; + if let Ok((stream, sa)) = self.listener.accept() { + info!("new connection from {}", sa); + self.listen = Some(stream); + } + if let Some(l) = &mut self.listen { + if let Ok(a) = l.read(&mut read_buf) { + if let Ok((r, _)) = + decode_from_slice::, _>(&read_buf[0..a], BINCODE_CONFIG) + { + return Some(r); + } + } + } + + return None; + } + } + + #[derive(Debug)] + pub struct UdpConnector { + socket: UdpSocket, + } + + impl UdpConnector { + pub fn new( + local_addr: &'static str, + local_port: u16, + remote_addr: &'static str, + remote_port: u16, + ) -> Result { + let socket = UdpSocket::bind( + (local_addr, local_port) + .to_socket_addrs()? + .next() + .ok_or(anyhow!("failed to parse ip address"))?, + )?; + socket.connect( + (remote_addr, remote_port) + .to_socket_addrs()? + .next() + .ok_or(anyhow!("failed to parse ip address"))?, + )?; + Ok(Self { socket }) + } + } + + impl Connector for UdpConnector { + fn send(&mut self, msgs: Vec) { + let r = encode_to_vec(&msgs, BINCODE_CONFIG).expect("failed to serialize udp packet"); + self.socket.send(&r).ok(); + } + + fn recv(&mut self) -> Option> { + let mut read_buf = [0_u8; 4096]; + if let Ok(a) = self.socket.recv(&mut read_buf) { + if let Ok((r, _)) = + decode_from_slice::, _>(&read_buf[0..a], BINCODE_CONFIG) + { + return Some(r); + } + } + + return None; + } + } +} + +#[cfg(feature = "std")] +pub use connector_std::*; diff --git a/rfe/src/lib.rs b/rfe/src/lib.rs new file mode 100644 index 0000000..8007abd --- /dev/null +++ b/rfe/src/lib.rs @@ -0,0 +1,30 @@ +#![no_std] + +#[cfg(feature = "std")] +extern crate std; + +pub mod connector; +use bincode::config::Configuration; +pub mod msg; + +mod rfe; +pub use rfe::*; + +#[cfg(feature = "to_csv")] +mod to_csv; +#[cfg(feature = "to_csv")] +pub use to_csv::*; + +pub mod time; +pub mod utils; + +pub const BINCODE_CONFIG: Configuration = bincode::config::standard(); + +#[macro_export] +macro_rules! unwrap_print_err { + ($x:expr, $msg: tt) => { + if let Err(_) = $x { + error!($msg) + } + }; +} diff --git a/rfe/src/msg.rs b/rfe/src/msg.rs new file mode 100644 index 0000000..7b0a12d --- /dev/null +++ b/rfe/src/msg.rs @@ -0,0 +1,222 @@ +extern crate alloc; + +use alloc::string::String; +use alloc::vec::Vec; +use bincode::{Decode, Encode}; +use macros::Kind; + +use crate::time::Timestamp; +#[cfg(feature = "to_csv")] +use crate::ToCsv; +#[cfg(feature = "to_csv")] +use macros::ToCsv; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct TargetMsg { + // instance is FROM for tlm, is TO for cmds + pub instance: Instance, + pub msg: MsgKind, +} + +impl TargetMsg { + pub fn new(instance: Instance, msg: MsgKind) -> Self { + Self { instance, msg } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub enum Instance { + All, + Other, + Example, + Example2, +} + +#[derive(Debug, Clone, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct ReinitAppCmd { + app_name: String, +} + +#[derive(Debug, Clone, PartialEq, Kind, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub enum Msg { + None, + SubRequest, + SubList(SubList), + SetTimeCmd(u64), + ReinitApp(ReinitAppCmd), + ExampleHk(ExampleHk), + ExampleOutData(ExampleOutData), + ExampleCmd(ExampleCmd), + DsHk(DsHk), + DsOutData(DsOutData), + DsCmd(DsCmd), + HsHk(HsHk), + HsOutData(HsOutData), + HsCmd(HsCmd), + ToHk(ToHk), + ToOutData(ToOutData), + ToCmd(ToCmd), +} + +#[derive(Debug, Clone, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct MsgPacket { + pub instance: Instance, + pub msg: Msg, + pub timestamp: Timestamp, +} + +impl MsgPacket { + pub fn to_target(&self) -> TargetMsg { + TargetMsg::new(self.instance, self.msg.kind()) + } + + pub fn new(instance: Instance, msg: Msg, timestamp: Timestamp) -> Self { + Self { + instance, + msg, + timestamp, + } + } + + pub fn timestamp(&mut self, timestamp: Timestamp) { + self.timestamp = timestamp; + } +} + +#[derive(Debug, Clone, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct SubList { + pub subs: Vec, +} + +#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct ExampleHk { + pub counter: u32, +} + +#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct ExampleOutData { + pub counter: u32, +} + +#[derive(Debug, Clone, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub enum ExampleCmd { + Noop, + Reset, +} + +#[derive(Debug, Clone, PartialEq, Eq, Hash, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct DsTlmSet { + pub items: Vec, + pub id: TlmSetId, + pub enabled: bool, + pub path: String, +} + +pub type TlmSetId = u16; + +#[derive(Debug, Clone, PartialEq, Eq, Hash, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct TlmSetItem { + pub target: TargetMsg, + pub decimation: u16, + pub counter: u16, +} + +#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] + +pub struct DsHk { + pub counter: u32, +} + +#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct DsOutData { + pub counter: u32, + pub bytes_written: u32, + pub bytes_written_this_cycle: u32, +} + +#[derive(Debug, Clone, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub enum DsCmd { + Noop, + Reset, + CloseAll, + Close(TlmSetId), + AddTlmSet(DsTlmSet), + RemoveTlmSet(TlmSetId), + DisableTlmSet(TlmSetId), + EnablTlmSet(TlmSetId), +} + +#[derive(Debug, Default, Clone, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct HsHk { + pub counter: u32, + pub cpu_usage: Vec, + pub mem_usage: u8, + pub fs_usage: Vec, + pub temps: Vec, + pub cmd_counter: u8, + pub cpu_usage_enabled: bool, + pub mem_usage_enabled: bool, + pub fs_usage_enabled: bool, +} + +#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct HsOutData { + pub counter: u32, +} + +#[derive(Debug, Clone, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub enum HsCmd { + Noop, + Reset, + WatchdogEnableManual(bool), + WatchdogEnableAuto(bool), + WatchdogResumeAuto, +} + +#[derive(Debug, Clone, PartialEq, Eq, Hash, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct ToTlmSet { + pub items: Vec, + pub id: TlmSetId, + pub enabled: bool, +} + +#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct ToHk { + pub counter: u32, +} + +#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub struct ToOutData { + pub counter: u32, +} + +#[derive(Debug, Clone, PartialEq, Encode, Decode)] +#[cfg_attr(feature = "to_csv", derive(ToCsv))] +pub enum ToCmd { + Noop, + Reset, + AddTlmSet(ToTlmSet), + RemoveTlmSet(TlmSetId), + DisableTlmSet(TlmSetId), + EnablTlmSet(TlmSetId), +} diff --git a/rfe/src/rfe.rs b/rfe/src/rfe.rs new file mode 100644 index 0000000..b8b05da --- /dev/null +++ b/rfe/src/rfe.rs @@ -0,0 +1,388 @@ +extern crate alloc; +use core::cell::RefCell; + +use alloc::rc::Rc; +use alloc::{collections::vec_deque::VecDeque, vec::Vec}; +use anyhow::{anyhow, Result}; +use hashbrown::{HashMap, HashSet}; +use log::*; + +use crate::{ + connector::Connector, + msg::{Instance, Msg, MsgPacket, SubList, TargetMsg}, + time::{TimeData, TimeDriver}, +}; + +pub trait Hk: Sized + Clone + Copy + 'static + Send + Sync {} +impl Hk for T where T: Sized + Clone + Copy + 'static + Send + Sync {} + +pub trait OutData: Sized + Clone + Copy + 'static + Send + Sync {} +impl OutData for T where T: Sized + Clone + Copy + 'static + Send + Sync {} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +pub enum Rate { + Hz1, + Hz5, + Hz10, + Hz20, + Hz50, + Hz100, +} + +type RfeTimeRef<'a> = Rc>>; + +pub struct RfeTime<'a> { + time_data: TimeData, + time_driver: &'a dyn TimeDriver, +} + +pub struct RfeInstance<'a> { + app_list: HashMap<&'a str, AppRef<'a>>, + time: RfeTimeRef<'a>, + #[allow(dead_code)] + instance: Instance, + connectors: Vec>, + sch_counter: u64, +} + +pub struct AppRef<'a> { + app: &'a mut dyn App, + app_rate: Rate, + out_data_rate: Rate, + hk_rate: Rate, + rfe: Rfe<'a>, +} + +pub struct Rfe<'a> { + subscriptions: HashSet, + msgs_to_send: Vec, + msgs_recevied: VecDeque, + instance: Instance, + time: RfeTimeRef<'a>, + subs_updated: bool, +} + +#[derive(Debug)] +pub struct ConnectorState<'a> { + connector: &'a mut dyn Connector, + subscriptions: HashSet, + subs_received: bool, + subs_last_requested: u64, +} + +impl<'a> Rfe<'a> { + pub fn new(instance: Instance, time: RfeTimeRef<'a>) -> Self { + Self { + subscriptions: HashSet::new(), + msgs_to_send: Vec::new(), + msgs_recevied: VecDeque::new(), + instance, + subs_updated: false, + time, + } + } + + pub fn get_instance(&self) -> Instance { + return self.instance; + } + + pub fn subscribe(&mut self, msg: TargetMsg) { + self.subscriptions.insert(msg); + self.subs_updated = true; + } + + pub fn subscribe_all>(&mut self, msgs: T) { + self.subscriptions.extend(msgs.into_iter()); + self.subs_updated = true; + } + + pub fn unsubscribe(&mut self, msg: &TargetMsg) { + self.subscriptions.remove(msg); + self.subs_updated = true; + } + + pub fn unsubscribe_all(&mut self) { + self.subscriptions.clear(); + self.subs_updated = true; + } + + pub fn send(&mut self, msg: Msg) { + self.msgs_to_send.push(MsgPacket::new( + self.get_instance(), + msg, + self.get_system_time(), + )); + } + + pub fn send_cmd(&mut self, msg: Msg, target: Instance) { + self.msgs_to_send + .push(MsgPacket::new(target, msg, self.get_system_time())); + } + + pub fn post_message(&mut self, msg: MsgPacket) { + self.msgs_recevied.push_back(msg); + } + + pub fn recv(&mut self) -> Option { + self.msgs_recevied.pop_front() + } + + /// time in 10ms increments starting from power on + pub fn get_met_time(&self) -> u64 { + self.time.borrow().time_data.sch_counter + } + + /// time in 10ms increments based on set time. Shows MET until time is set + pub fn get_time(&self) -> u64 { + let time = self.time.borrow(); + time.time_data.sch_counter + time.time_data.time_offset + } + + /// Time in microseconds relative to system epoch + pub fn get_system_time(&self) -> u64 { + let time = self.time.borrow(); + time.time_driver.get_system_time(time.time_data) + } +} + +pub trait App { + fn init(&mut self, rfe: &mut Rfe) -> Result<()>; + fn run(&mut self, rfe: &mut Rfe); + fn hk(&mut self, rfe: &mut Rfe); + fn out_data(&mut self, rfe: &mut Rfe); + fn get_app_rate(&self) -> Rate; +} + +impl<'a> RfeInstance<'a> { + pub fn new(instance: Instance, time_driver: &'a dyn TimeDriver) -> Self { + let time = Rc::new(RefCell::new(RfeTime { + time_data: TimeData { + sch_counter: 0, + time_offset: 0, + }, + time_driver, + })); + Self { + app_list: HashMap::new(), + instance, + connectors: Vec::new(), + time, + sch_counter: 0, + } + } + + pub fn add_app(&mut self, name: &'a str, app: &'a mut dyn App) -> Result<()> { + if self.app_list.contains_key(name) { + return Err(anyhow!( + "failed to add app {name}, already added an app with that name" + )); + } + let app_rate = app.get_app_rate(); + self.app_list.insert( + name, + AppRef { + app: app, + app_rate: app_rate, + hk_rate: Rate::Hz1, + out_data_rate: app_rate, + rfe: Rfe::new(self.instance, self.time.clone()), + }, + ); + + let appref = self.app_list.get_mut(name).unwrap(); + if let Err(e) = appref.app.init(&mut appref.rfe) { + error!("app {name} failed to initialize {e}"); + } + + return Ok(()); + } + + pub fn add_connector(&mut self, connector: &'a mut dyn Connector) { + self.connectors.push(ConnectorState { + connector, + subs_last_requested: 0, + subs_received: false, + subscriptions: HashSet::new(), + }); + } + + /// Expected to be called at 100Hz + pub fn run(&mut self) { + let mut msgs = Vec::new(); + for app in self.app_list.values_mut() { + if app.app_rate == Rate::Hz100 + || (self.sch_counter % 2 == 0 && app.app_rate == Rate::Hz50) + || (self.sch_counter % 5 == 0 && app.app_rate == Rate::Hz20) + || (self.sch_counter % 10 == 0 && app.app_rate == Rate::Hz10) + || (self.sch_counter % 20 == 0 && app.app_rate == Rate::Hz5) + || (self.sch_counter % 100 == 0 && app.app_rate == Rate::Hz1) + { + app.app.run(&mut app.rfe); + } + + if app.hk_rate == Rate::Hz100 + || (self.sch_counter % 2 == 0 && app.hk_rate == Rate::Hz50) + || (self.sch_counter % 5 == 0 && app.hk_rate == Rate::Hz20) + || (self.sch_counter % 10 == 0 && app.hk_rate == Rate::Hz10) + || (self.sch_counter % 20 == 0 && app.hk_rate == Rate::Hz5) + || (self.sch_counter % 100 == 0 && app.hk_rate == Rate::Hz1) + { + app.app.hk(&mut app.rfe); + } + + if app.out_data_rate == Rate::Hz100 + || (self.sch_counter % 2 == 0 && app.out_data_rate == Rate::Hz50) + || (self.sch_counter % 5 == 0 && app.out_data_rate == Rate::Hz20) + || (self.sch_counter % 10 == 0 && app.out_data_rate == Rate::Hz10) + || (self.sch_counter % 20 == 0 && app.out_data_rate == Rate::Hz5) + || (self.sch_counter % 100 == 0 && app.out_data_rate == Rate::Hz1) + { + app.app.out_data(&mut app.rfe); + } + + let new_msgs = core::mem::take(&mut app.rfe.msgs_to_send); + msgs.extend(new_msgs); + } + + for app in self.app_list.values_mut() { + for msg in &msgs { + if app + .rfe + .subscriptions + .contains(&TargetMsg::new(self.instance, msg.msg.kind())) + || app + .rfe + .subscriptions + .contains(&TargetMsg::new(Instance::All, msg.msg.kind())) + { + app.rfe.post_message(msg.clone()); + } + } + } + + // send message to connectors + for connector_state in &mut self.connectors { + let mut to_send = Vec::new(); + for msg in &msgs { + if connector_state + .subscriptions + .contains(&TargetMsg::new(self.instance, msg.msg.kind())) + || connector_state + .subscriptions + .contains(&TargetMsg::new(Instance::All, msg.msg.kind())) + || connector_state + .subscriptions + .contains(&TargetMsg::new(Instance::Other, msg.msg.kind())) + { + to_send.push(msg.clone()); + } + } + if to_send.len() > 0 { + connector_state.connector.send(to_send); + } + } + + drop(msgs); + let mut connector_msgs = Vec::new(); + + // receive messages from connectors + for connector_state in &mut self.connectors { + if let Some(msgs) = connector_state.connector.recv() { + // check for sublist/sub request + for msg in &msgs { + if let Msg::SubList(list) = &msg.msg { + connector_state.subs_received = true; + connector_state.subscriptions.clear(); + connector_state.subscriptions.extend(list.subs.clone()); + } + if let Msg::SetTimeCmd(new_time) = &msg.msg { + self.time.borrow_mut().time_data.time_offset = *new_time; + } + if let Msg::SubRequest = msg.msg { + let mut subs = Vec::new(); + for app in self.app_list.values() { + subs.extend(app.rfe.subscriptions.clone()); + } + let mut sublist = Vec::new(); + sublist.push(MsgPacket { + instance: self.instance, + msg: Msg::SubList(SubList { subs }), + timestamp: 0, + }); + connector_state.connector.send(sublist); + } + } + + connector_msgs.extend(msgs); + } + } + + // send messages to apps + for app in self.app_list.values_mut() { + for msg in &connector_msgs { + if app + .rfe + .subscriptions + .contains(&TargetMsg::new(msg.instance, msg.msg.kind())) + || app + .rfe + .subscriptions + .contains(&TargetMsg::new(Instance::All, msg.msg.kind())) + || app + .rfe + .subscriptions + .contains(&TargetMsg::new(Instance::Other, msg.msg.kind())) + { + app.rfe.post_message(msg.clone()); + } + } + } + + // clear connector subs if subs changed + for rfe in self.app_list.values_mut() { + if rfe.rfe.subs_updated { + for connector_state in &mut self.connectors { + connector_state.subs_received = false; + } + rfe.rfe.subs_updated = false; + } + } + + // handle connector subscriptions + for connector_state in &mut self.connectors { + if (!connector_state.subs_received + && self.sch_counter - connector_state.subs_last_requested >= 10) + || (connector_state.subs_received + && self.sch_counter - connector_state.subs_last_requested >= 10000) + { + // request subs + let mut request = Vec::new(); + request.push(MsgPacket { + instance: self.instance, + msg: Msg::SubRequest, + timestamp: 0, + }); + connector_state.connector.send(request); + connector_state.subs_last_requested = self.sch_counter; + } + } + + self.sch_counter += 1; + } + + #[cfg(feature = "std")] + pub fn start(&mut self) { + use core::time::Duration; + use std::{thread::sleep, time::Instant}; + + let mut next_time = Instant::now() + Duration::from_millis(10); + loop { + sleep(next_time - Instant::now()); + self.run(); + next_time += Duration::from_millis(10); + if next_time < Instant::now() { + next_time = Instant::now() + Duration::from_millis(10); + } + } + } +} diff --git a/rfe/src/time.rs b/rfe/src/time.rs new file mode 100644 index 0000000..ea6d6ce --- /dev/null +++ b/rfe/src/time.rs @@ -0,0 +1,45 @@ +pub type Timestamp = u64; + +#[derive(Debug, Clone, Copy)] +pub struct TimeData { + pub sch_counter: u64, + pub time_offset: u64, +} + +pub trait TimeDriver { + /// Time in microseconds relative to system epoch + fn get_system_time(&self, time_data: TimeData) -> Timestamp; +} + +#[cfg(feature = "std")] +pub struct UnixTimeDriver; +#[cfg(feature = "std")] +impl UnixTimeDriver { + pub fn new() -> Self { + Self {} + } +} + +#[cfg(feature = "std")] +impl TimeDriver for UnixTimeDriver { + fn get_system_time(&self, _time_data: TimeData) -> Timestamp { + use std::time::SystemTime; + SystemTime::now() + .duration_since(SystemTime::UNIX_EPOCH) + .unwrap() + .as_micros() as u64 + } +} + +pub struct SchTimeDriver; +impl SchTimeDriver { + pub fn new() -> Self { + Self {} + } +} + +impl TimeDriver for SchTimeDriver { + fn get_system_time(&self, time_data: TimeData) -> Timestamp { + time_data.sch_counter + time_data.time_offset + } +} diff --git a/rfe/src/to_csv.rs b/rfe/src/to_csv.rs new file mode 100644 index 0000000..52b7088 --- /dev/null +++ b/rfe/src/to_csv.rs @@ -0,0 +1,151 @@ +extern crate alloc; +use alloc::format; +use alloc::string::String; +use alloc::vec; +use alloc::vec::Vec; + +pub trait ToCsv { + fn to_csv(&self) -> Vec; +} + +pub trait ToCsvClean { + fn to_csv_clean(&self) -> Vec; +} + +impl ToCsv for u8 { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for i8 { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for u16 { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for i16 { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for u32 { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for i32 { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for u64 { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for i64 { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for &str { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for String { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for bool { + fn to_csv(&self) -> Vec { + vec![format!(" = {}", self)] + } +} + +impl ToCsv for Vec { + fn to_csv(&self) -> Vec { + self.iter() + .enumerate() + .map(|(i, x)| { + x.to_csv() + .iter() + .map(|a| format!("[{}].{}", i, a)) + .collect::>() + }) + .flatten() + .collect() + } +} + +impl ToCsv for [T; N] { + fn to_csv(&self) -> Vec { + self.iter() + .enumerate() + .map(|(i, x)| { + x.to_csv() + .iter() + .map(|a| format!("[{}].{}", i, a)) + .collect::>() + }) + .flatten() + .collect() + } +} + +impl ToCsv for (T1, T2) { + fn to_csv(&self) -> Vec { + let mut csvs = Vec::new(); + csvs.extend(self.0.to_csv().iter().map(|x| format!("[0].{}", x))); + csvs.extend(self.1.to_csv().iter().map(|x| format!("[1].{}", x))); + csvs + } +} + +impl ToCsv for (T1, T2, T3) { + fn to_csv(&self) -> Vec { + let mut csvs = Vec::new(); + csvs.extend(self.0.to_csv().iter().map(|x| format!("[0].{}", x))); + csvs.extend(self.1.to_csv().iter().map(|x| format!("[1].{}", x))); + csvs.extend(self.2.to_csv().iter().map(|x| format!("[2].{}", x))); + csvs + } +} + +impl ToCsv for (T1, T2, T3, T4) { + fn to_csv(&self) -> Vec { + let mut csvs = Vec::new(); + csvs.extend(self.0.to_csv().iter().map(|x| format!("[0].{}", x))); + csvs.extend(self.1.to_csv().iter().map(|x| format!("[1].{}", x))); + csvs.extend(self.2.to_csv().iter().map(|x| format!("[2].{}", x))); + csvs.extend(self.3.to_csv().iter().map(|x| format!("[3].{}", x))); + csvs + } +} + +impl ToCsvClean for T { + fn to_csv_clean(&self) -> Vec { + self.to_csv().iter().map(|x| to_csv_clean(x)).collect() + } +} + +fn to_csv_clean(s: &String) -> String { + let s = s.replace(".[", "["); + let s = s.replace(". =", " ="); + s +} diff --git a/rfe/src/utils.rs b/rfe/src/utils.rs new file mode 100644 index 0000000..a11a999 --- /dev/null +++ b/rfe/src/utils.rs @@ -0,0 +1,63 @@ + +pub struct ManualAuto { + value_auto: T, + value_manual: T, + is_manual: bool, + has_changed: bool, +} + +impl ManualAuto { + pub fn new(value: T, is_manual: bool) -> Self { + Self { + value_auto: value.clone(), + value_manual: value, + is_manual, + has_changed: false, + } + } + + pub fn auto_set(&mut self, value: T) { + if self.value_auto != value && !self.is_manual { + self.has_changed = true; + } + self.value_auto = value; + } + + pub fn manual_set(&mut self, value: T) { + if (self.value_manual != value && self.is_manual) + || (!self.is_manual && self.value_auto != value) + { + self.has_changed = true; + } + self.is_manual = true; + self.value_manual = value; + } + + pub fn resume_auto(&mut self) { + if self.is_manual && self.value_manual != self.value_auto { + self.has_changed = true; + } + self.is_manual = false; + } + + pub fn get(&self) -> &T { + if self.is_manual { + &self.value_manual + } else { + &self.value_auto + } + } + + pub fn is_manual(&self) -> bool { + return self.is_manual; + } + + /// Checks if the internal value has changed since the last call to has_changed + pub fn has_changed(&mut self) -> bool { + if self.has_changed { + self.has_changed = false; + return true; + } + return false; + } +} diff --git a/tools/decom/Cargo.toml b/tools/decom/Cargo.toml new file mode 100644 index 0000000..0a36406 --- /dev/null +++ b/tools/decom/Cargo.toml @@ -0,0 +1,15 @@ +[package] +name = "decom" +version = "0.1.0" +edition = "2021" + +[[bin]] +name = "decom" +path = "src/decom.rs" + +[dependencies] +anyhow.workspace = true +log.workspace = true +simple_logger.workspace = true +rfe = { path = "../../rfe", features = ["to_csv"] } +bincode = { workspace = true, features = ["std"] } diff --git a/tools/decom/src/decom.rs b/tools/decom/src/decom.rs new file mode 100644 index 0000000..16bf51a --- /dev/null +++ b/tools/decom/src/decom.rs @@ -0,0 +1,113 @@ +#![feature(bufreader_peek)] +use std::{ + collections::HashMap, + env::args, + fs::{read_dir, OpenOptions}, + io::{BufReader, Write}, + path::PathBuf, + str::FromStr, + thread::spawn, +}; + +use anyhow::Result; +use bincode::decode_from_std_read; +use log::*; +use rfe::ToCsvClean; +use rfe::{msg::MsgPacket, BINCODE_CONFIG}; +use simple_logger::SimpleLogger; + +fn decom_file(file_path: String, out_dir: String) -> Result<()> { + info!("decomming {file_path}"); + let f = OpenOptions::new().read(true).open(file_path)?; + let mut rea = BufReader::new(f); + let mut files = HashMap::new(); + + 'outer: loop { + while let Ok(_) = rea.peek(1) { + let msg = match decode_from_std_read::(&mut rea, BINCODE_CONFIG) { + Ok(m) => m, + Err(_) => break 'outer, + }; + if !files.contains_key(&msg.msg.kind()) { + let path = PathBuf::from_str(&out_dir)? + .join(format!("{:?}.csv", msg.msg.kind()).to_lowercase()); + let exists = path.exists(); + let mut w = OpenOptions::new() + .write(true) + .create(true) + .append(true) + .open(path)?; + + if !exists { + let values = msg.to_csv_clean(); + let mut line = values + .iter() + .map(|x| x.split_once(" = ").expect("unexpected csv data").0) + .collect::>() + .join(","); + line += "\n"; + w.write(line.as_bytes())?; + } + + files.insert(msg.msg.kind(), w); + } + + if let Some(w) = files.get_mut(&msg.msg.kind()) { + let values = msg.to_csv_clean(); + let mut line = values + .iter() + .map(|x| x.split_once(" = ").expect("unexpected csv data").1) + .collect::>() + .join(","); + line += "\n"; + w.write(line.as_bytes())?; + } + } + } + return Ok(()); +} + +fn decom_task(folder: String, out_dir: String) -> Result<()> { + let mut tasks = Vec::new(); + let mut files = Vec::new(); + for ent in read_dir(folder)? { + let ent = ent?; + + if ent.metadata()?.is_dir() { + let folder = ent.path().to_string_lossy().to_string(); + let out_dir = out_dir.clone(); + tasks.push(spawn(move || decom_task(folder, out_dir))); + } else if ent.metadata()?.is_file() { + let file_path = ent.path().to_string_lossy().to_string(); + files.push(file_path); + } + } + + files.sort(); + for file_path in files.iter() { + let out_dir = out_dir.clone(); + if let Err(e) = decom_file(file_path.clone(), out_dir) { + error!("{e}"); + } + } + + for task in tasks { + task.join().ok(); + } + + return Ok(()); +} + +fn main() -> Result<()> { + SimpleLogger::new().init().unwrap(); + info!("started decom"); + + let args = args().skip(1).collect::>(); + let folder = args[0].clone(); + let out_dir = args[1].clone(); + + decom_task(folder, out_dir)?; + + info!("finished decom"); + return Ok(()); +}