This commit is contained in:
Jacob Weaver
2024-12-13 10:23:45 -06:00
commit 6b27e84e4b
38 changed files with 2821 additions and 0 deletions
View File
+21
View File
@@ -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/
+3
View File
@@ -0,0 +1,3 @@
{
"editor.formatOnSave": true
}
+24
View File
@@ -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"
+19
View File
@@ -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
+192
View File
@@ -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<TlmSetId, DsFile>,
enabled: bool,
}
pub struct DsFileSettings {
pub max_size: u32,
pub max_age: u32,
pub enabled: bool,
}
pub struct Ds {
data: DsData,
tlm_sets: HashMap<TlmSetId, DsTlmSet>,
start_enabled: bool,
}
impl Ds {
pub fn new(tlm_sets: HashMap<TlmSetId, DsTlmSet>, 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
}
}
+72
View File
@@ -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<File>,
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<usize> {
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(());
}
}
}
+25
View File
@@ -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<usize> {
return Ok(0);
}
pub fn flush(&mut self) -> Result<()> {
return Ok(());
}
}
+12
View File
@@ -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
+69
View File
@@ -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
}
}
+22
View File
@@ -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",
] }
+154
View File
@@ -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<u8>;
fn check_mem_usage(&mut self) -> u8;
fn check_fs_usage(&mut self) -> Vec<u8>;
fn check_temps(&mut self) -> Vec<i8>;
}
#[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<bool>,
}
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
}
}
+59
View File
@@ -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<u8> {
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<u8> {
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<i8> {
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()
}
}
+84
View File
@@ -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<Self> {
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();
}
}
}
+13
View File
@@ -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
+151
View File
@@ -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<TlmSetId, ToTlmSet>,
}
impl<'a> To<'a> {
pub fn new(connector: &'a mut dyn Connector, tlm_sets: HashMap<TlmSetId, ToTlmSet>) -> 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
}
}
+15
View File
@@ -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
+140
View File
@@ -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(());
}
+10
View File
@@ -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
+30
View File
@@ -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);
}
}
@@ -0,0 +1,6 @@
[target.'cfg(all(target_arch = "arm", target_os = "none"))']
runner = "elf2uf2-rs -d"
[build]
target = "thumbv6m-none-eabi"
+39
View File
@@ -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"
+25
View File
@@ -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");
}
+3
View File
@@ -0,0 +1,3 @@
#!/bin/bash
cargo run --release --target thumbv6m-none-eabi
+15
View File
@@ -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;
+204
View File
@@ -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<MsgPacket>) {
info!("got msgs: {:?}", msgs);
}
fn recv(&mut self) -> Option<Vec<MsgPacket>> {
return None;
}
}
struct BlinkApp<'a> {
led_pin: &'a mut Pin<Gpio25, FunctionSio<SioOutput>, 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<Gpio25, FunctionSio<SioOutput>, PullNone>,
usb: Option<Driver<'static, USB>>,
}
struct UsbBinding;
unsafe impl typelevel::Binding<typelevel::USBCTRL_IRQ, usb::InterruptHandler<USB>> 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<u8>; 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::<PullNone>()
.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 {
<usb::InterruptHandler<USB> as typelevel::Handler<typelevel::USBCTRL_IRQ>>::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::<u64, 1, 1000000>::from_ticks(10000);
loop {
instance.run();
Mono::delay_until(next_time).await;
next_time += Duration::<u64, 1, 1000000>::from_ticks(10000);
}
}
}
+17
View File
@@ -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"
+14
View File
@@ -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"] }
+169
View File
@@ -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<Self> {
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<String> = (&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<String> {
extern crate alloc;
use alloc::format;
let mut values:Vec<String> = 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<String> = (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<String> {
extern crate alloc;
use alloc::format;
let mut values: Vec<String> = 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)
}
+187
View File
@@ -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<MsgPacket>);
fn recv(&mut self) -> Option<Vec<MsgPacket>>;
}
#[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<Vec<MsgPacket>>,
receiver: Receiver<Vec<MsgPacket>>,
}
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<MsgPacket>) {
self.sender.send(msgs).ok();
}
fn recv(&mut self) -> Option<Vec<MsgPacket>> {
self.receiver.recv_timeout(Duration::ZERO).ok()
}
}
#[derive(Debug)]
pub struct TcpConnector {
listener: TcpListener,
listen: Option<TcpStream>,
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<Self> {
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<MsgPacket>) {
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<Vec<MsgPacket>> {
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::<Vec<MsgPacket>, _>(&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<Self> {
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<MsgPacket>) {
let r = encode_to_vec(&msgs, BINCODE_CONFIG).expect("failed to serialize udp packet");
self.socket.send(&r).ok();
}
fn recv(&mut self) -> Option<Vec<MsgPacket>> {
let mut read_buf = [0_u8; 4096];
if let Ok(a) = self.socket.recv(&mut read_buf) {
if let Ok((r, _)) =
decode_from_slice::<Vec<MsgPacket>, _>(&read_buf[0..a], BINCODE_CONFIG)
{
return Some(r);
}
}
return None;
}
}
}
#[cfg(feature = "std")]
pub use connector_std::*;
+30
View File
@@ -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)
}
};
}
+222
View File
@@ -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<TargetMsg>,
}
#[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<TlmSetItem>,
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<u8>,
pub mem_usage: u8,
pub fs_usage: Vec<u8>,
pub temps: Vec<i8>,
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<TlmSetItem>,
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),
}
+388
View File
@@ -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<T> Hk for T where T: Sized + Clone + Copy + 'static + Send + Sync {}
pub trait OutData: Sized + Clone + Copy + 'static + Send + Sync {}
impl<T> 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<RefCell<RfeTime<'a>>>;
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<ConnectorState<'a>>,
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<TargetMsg>,
msgs_to_send: Vec<MsgPacket>,
msgs_recevied: VecDeque<MsgPacket>,
instance: Instance,
time: RfeTimeRef<'a>,
subs_updated: bool,
}
#[derive(Debug)]
pub struct ConnectorState<'a> {
connector: &'a mut dyn Connector,
subscriptions: HashSet<TargetMsg>,
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<T: IntoIterator<Item = TargetMsg>>(&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<MsgPacket> {
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);
}
}
}
}
+45
View File
@@ -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
}
}
+151
View File
@@ -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<String>;
}
pub trait ToCsvClean {
fn to_csv_clean(&self) -> Vec<String>;
}
impl ToCsv for u8 {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for i8 {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for u16 {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for i16 {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for u32 {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for i32 {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for u64 {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for i64 {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for &str {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for String {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl ToCsv for bool {
fn to_csv(&self) -> Vec<String> {
vec![format!(" = {}", self)]
}
}
impl<T: ToCsv> ToCsv for Vec<T> {
fn to_csv(&self) -> Vec<String> {
self.iter()
.enumerate()
.map(|(i, x)| {
x.to_csv()
.iter()
.map(|a| format!("[{}].{}", i, a))
.collect::<Vec<String>>()
})
.flatten()
.collect()
}
}
impl<T: ToCsv, const N: usize> ToCsv for [T; N] {
fn to_csv(&self) -> Vec<String> {
self.iter()
.enumerate()
.map(|(i, x)| {
x.to_csv()
.iter()
.map(|a| format!("[{}].{}", i, a))
.collect::<Vec<String>>()
})
.flatten()
.collect()
}
}
impl<T1: ToCsv, T2: ToCsv> ToCsv for (T1, T2) {
fn to_csv(&self) -> Vec<String> {
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<T1: ToCsv, T2: ToCsv, T3: ToCsv> ToCsv for (T1, T2, T3) {
fn to_csv(&self) -> Vec<String> {
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<T1: ToCsv, T2: ToCsv, T3: ToCsv, T4: ToCsv> ToCsv for (T1, T2, T3, T4) {
fn to_csv(&self) -> Vec<String> {
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<T: ToCsv> ToCsvClean for T {
fn to_csv_clean(&self) -> Vec<String> {
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
}
+63
View File
@@ -0,0 +1,63 @@
pub struct ManualAuto<T: Clone + PartialEq> {
value_auto: T,
value_manual: T,
is_manual: bool,
has_changed: bool,
}
impl<T: Clone + PartialEq> ManualAuto<T> {
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;
}
}
+15
View File
@@ -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"] }
+113
View File
@@ -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::<MsgPacket, _, _>(&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::<Vec<&str>>()
.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::<Vec<&str>>()
.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::<Vec<String>>();
let folder = args[0].clone();
let out_dir = args[1].clone();
decom_task(folder, out_dir)?;
info!("finished decom");
return Ok(());
}