Compare commits

...
11 Commits
Author SHA1 Message Date
jcweaver997 d7d8c0789a changed decom to output json 2025-04-07 22:32:22 -05:00
jcweaver997 bfe3333d65 docs 2025-04-02 22:50:55 -05:00
Jacob Weaver e645bbcc08 ui stuff 2024-12-19 12:12:55 -06:00
Jacob Weaver 70b2f1386d made serial config pub 2024-12-19 10:10:12 -06:00
Jacob Weaver 27559a8d31 added serial 2024-12-19 09:52:40 -06:00
jcweaver997 f59b164095 reflect 2024-12-18 22:13:41 -06:00
Jacob Weaver e9a4f3de75 reflect do vec 2024-12-18 16:16:26 -06:00
Jacob Weaver fbd47eba43 reflect do vec 2024-12-18 16:15:55 -06:00
Jacob Weaver 1a511523f4 reflect 2024-12-18 15:44:02 -06:00
Jacob Weaver c9ad723e79 start reflect 2024-12-17 16:13:03 -06:00
Jacob Weaver 5401d624bd added example app to rp pico 2024-12-17 11:12:30 -06:00
31 changed files with 1912 additions and 403 deletions
+5
View File
@@ -9,6 +9,9 @@ members = [
"builds/ground",
"builds/rp_pico_example",
"rfe",
"rfe/macros",
"rfe/reflect",
"rfe/reflect/reflect_macros",
"tools/decom",
]
resolver = "2"
@@ -24,3 +27,5 @@ bincode = { version = "2.0.0-rc.3", default-features = false, features = [
simple_logger = "5.0.0"
rp2040-hal = "0.10.2"
rp2040-pac = "0.6.0"
mio-serial = "=5.0.5"
reflect = { path = "rfe/reflect" }
+177
View File
@@ -0,0 +1,177 @@
# Rust Flight Executive (RFE)
RFE is a framework for building real-time embedded applications in Rust, inspired by concepts from NASA's Core Flight Executive (cFE). Rather than being a direct port of cFE to Rust, RFE is designed as an alternative that embraces Rust's idioms and strengths. Unlike cFE which primarily targets computer systems, RFE extends support to microcontrollers and other resource-constrained embedded devices. It provides a message-passing architecture for inter-application communication, time management, and scheduling at different rates. I/O operations are handled either by the standard library (when the `std` feature is enabled) or by board support packages like `rp2040-hal` for embedded platforms. These board support packages can be used generically through the traits defined in the `embedded-hal` crate, allowing for portable code across different embedded platforms.
## Overview
RFE is a framework that enables the development of modular, reusable flight software applications in Rust. While it shares some architectural concepts with NASA's cFE, RFE makes different design choices that are more ergonomic and better aligned with Rust's safety features, performance characteristics, and modern development ecosystem. By supporting both standard computers and microcontrollers, RFE offers greater flexibility for embedded systems development across a wide range of hardware platforms.
## Features
- **Message-passing architecture** for inter-application communication
- **Time management** for both system and monotonic time
- **Scheduling of applications** at different rates (1Hz to 100Hz)
- **Platform abstraction** through feature flags:
- `std`: Standard library support (Unix, Windows)
- `rp2040`: Raspberry Pi Pico support
- More platform abstractions planned for future releases
- Support for custom platforms by implementing a few simple traits
- **Runtime reflection**:
- `reflect`: Runtime type information that can be used for dynamic command generation and telemetry processing
- **Connectors**: A trait-based system for communication between instances, with built-in implementations (TCP, UDP, Memory) and support for custom user-defined connectors
- **Core applications**:
- Data Storage (DS): For telemetry storage
- Health & Safety (HS): For system monitoring
- Telemetry Output (TO): For telemetry management
- Example app: Demonstrating application development
## Architecture
RFE follows a modular architecture where applications communicate through a message bus. The core components include:
- **RfeInstance**: The main framework instance that manages applications and scheduling. An instance is intended as a runner for apps - multiple apps can be added to a single instance, but thread/task priority is given to each instance rather than to individual apps. For each app to have its own task priority, you will need a separate instance for each app.
- **App trait**: Interface that all applications must implement
- **Rfe**: Interface for applications to interact with the framework
- **Connectors**: A trait-based interface for communication between different RFE instances, allowing users to implement custom message transport mechanisms. Instances can communicate with each other through these connectors.
## Getting Started
### Prerequisites
- Rust toolchain (latest stable recommended)
- For RP2040 development: `thumbv6m-none-eabi` target
### Building
```bash
# For standard desktop build
cargo build --package example_build
# For Raspberry Pi Pico
cargo build --package rp_pico_example --target thumbv6m-none-eabi
```
### Creating an Application
Applications in RFE must implement the `App` trait:
```rust
use rfe::*;
use anyhow::Result;
struct MyApp;
impl App for MyApp {
fn init(&mut self, rfe: &mut Rfe) -> Result<()> {
// Initialize the application
// Set up message subscriptions
rfe.subscribe(TargetMsg::new(rfe.get_instance(), MsgKind::YourMsgKind));
Ok(())
}
fn run(&mut self, rfe: &mut Rfe) {
// Run the application logic
// Process received messages
while let Some(msg) = rfe.recv() {
// Handle message
}
}
fn hk(&mut self, rfe: &mut Rfe) {
// Generate housekeeping data
rfe.send(Msg::YourHkMsg(self.your_hk_data.clone()));
}
fn out_data(&mut self, rfe: &mut Rfe) {
// Generate output data
rfe.send(Msg::YourOutDataMsg(self.your_out_data.clone()));
}
fn get_app_rate(&self) -> Rate {
// Set the intended rate at which your app runs
Rate::Hz10 // Run at 10Hz
}
}
```
### Running an RFE Instance
```rust
use rfe::*;
use rfe::connector::MemConnector;
use std::thread;
use std::sync::Arc;
// Create a shared memory connector for inter-instance communication
let (connector1, connector2) = Arc::new(MemConnector::new());
// First instance (you can set the priority of the thread you run the instance in)
thread::spawn(move || {
// Create a time driver for the first instance
let time_driver = StdTimeDriver::new();
// Create the first RFE instance
let mut instance1 = RfeInstance::new(Instance::Main, &time_driver);
// Add the memory connector to the instance
instance1.add_connector(connector1);
// Add applications to the first instance
let mut my_app1 = MyApp::new();
instance1.add_app("my_app1", &mut my_app1).unwrap();
// Run the first instance
instance1.start();
});
// Second instance
let time_driver2 = StdTimeDriver::new();
let mut instance2 = RfeInstance::new(Instance::Secondary, &time_driver2);
// Add the memory connector to the second instance
instance2.add_connector(connector_clone);
// Add applications to the second instance
let mut my_app2 = AnotherApp::new();
instance2.add_app("my_app2", &mut my_app2).unwrap();
// Run the second instance
instance2.start();
```
In this example:
- Two RFE instances are created, each with its own application
- The first instance runs in a separate thread. You can set the priority of this thread
- The second instance runs in the main thread
- A shared `MemConnector` allows the instances to communicate with each other
- Each app can send messages that will be routed to the appropriate instance
## Project Structure
- `rfe/`: Core framework implementation
- `apps/`: Application implementations
- `ds/`: Data Storage application
- `hs/`: Health & Safety application
- `to/`: Telemetry Output application
- `example/`: Example application
- `builds/`: Example builds for different platforms
- `example_build/`: Standard desktop build
- `rp_pico_example/`: Raspberry Pi Pico build
- `ground/`: Ground system with GUI
- `tools/`: Development tools
- `decom/`: Telemetry decommutation tool
## Platform Support
RFE supports multiple platforms through feature flags:
- **Standard Desktop** (`std` feature): Unix, Windows
- **Raspberry Pi Pico** (`rp2040` feature): Embedded ARM Cortex-M0+
- Users can add support for any platform by implementing just a few key traits (primarily the `TimeDriver` trait)
- Board support packages utilize the `embedded-hal` traits for generic, portable hardware abstraction
- Platform-specific functionality can be easily integrated through Rust's trait system
- Additional platform support is planned for future releases
The `reflect` feature can be enabled on any platform to provide runtime type information.
+2
View File
@@ -41,7 +41,9 @@ impl App for Example {
ExampleCmd::Noop => info!("NOOP command received"),
ExampleCmd::Reset => {
info!("RESET command received");
let perf = self.data.hk.perf;
self.data = Default::default();
self.data.hk.perf = perf;
}
},
_ => {
+11 -4
View File
@@ -89,14 +89,16 @@ mod watchdog_rp2040 {
pub struct Rp2040Watchdog {
wd: Watchdog,
time: u32,
timeout: u32,
enabled: bool,
}
impl Rp2040Watchdog {
pub fn new(wd: Watchdog) -> Self {
Self {
wd,
time: 3 * 1000000,
timeout: 3 * 1000000,
enabled: false,
}
}
}
@@ -105,16 +107,21 @@ mod watchdog_rp2040 {
fn enable(&mut self) {
self.wd
.start(rp2040_hal::fugit::Duration::<u32, 1, 1000000>::from_ticks(
self.time,
self.timeout,
));
self.enabled = true;
}
fn disable(&mut self) {
self.wd.disable();
self.enabled = false;
}
fn set_timeout(&mut self, time: i32) {
self.time = time as u32 * 1000000;
self.timeout = time as u32 * 1000000;
if self.enabled {
self.enable();
}
}
fn feed(&mut self) {
+3 -1
View File
@@ -5,7 +5,9 @@ edition = "2021"
[dependencies]
rfe.path = "../../rfe"
rfe.features = ["std"]
rfe.features = ["std", "reflect"]
anyhow.workspace = true
simple_logger = "5.0.0"
log.workspace = true
egui = "0.30.0"
eframe = "0.30.0"
+193 -6
View File
@@ -1,24 +1,47 @@
#![feature(thread_sleep_until)]
#![cfg_attr(not(debug_assertions), windows_subsystem = "windows")] // hide console window on Windows in release
use std::{
thread::sleep_until,
collections::HashMap,
thread::{sleep, spawn},
time::{Duration, Instant},
};
use anyhow::anyhow;
use anyhow::Result;
use connector::{Connector, UdpConnector};
use egui::{Layout, Ui};
use log::*;
use msg::Msg;
use rfe::reflect::*;
use rfe::*;
extern crate alloc;
use simple_logger::SimpleLogger;
fn main() -> Result<()> {
SimpleLogger::new().init().unwrap();
#[derive(Debug, Default, Reflect)]
struct TestStruct {
counter: i8,
}
let mut udp = UdpConnector::new("127.0.0.1", 7011, "127.0.0.1", 7010)?;
#[derive(Debug, Default, Reflect)]
enum EnumTest {
#[default]
None,
Test(TestStruct),
}
fn main() -> Result<()> {
SimpleLogger::new()
.with_level(LevelFilter::Debug)
.init()
.unwrap();
spawn(|| {
let mut udp = UdpConnector::new("127.0.0.1", 7011, "127.0.0.1", 7010).unwrap();
let mut next_time = Instant::now() + Duration::from_millis(10);
loop {
sleep_until(next_time);
sleep(next_time - Instant::now());
while let Some(msgs) = udp.recv() {
for msg in msgs {
@@ -27,4 +50,168 @@ fn main() -> Result<()> {
}
next_time += Duration::from_millis(10);
}
});
let options = eframe::NativeOptions {
viewport: egui::ViewportBuilder::default().with_inner_size([560.0, 480.0]),
..Default::default()
};
eframe::run_native(
"RFE Ground",
options,
Box::new(|_cc| Ok(Box::<MyApp>::default())),
)
.or(Err(anyhow!("eframe failed")))?;
return Ok(());
}
struct MyApp {
command: Msg,
values: HashMap<String, String>,
}
impl Default for MyApp {
fn default() -> Self {
Self {
command: Msg::None,
values: HashMap::new(),
}
}
}
fn build_command_ui_value_number(ui: &mut Ui, name: &str, cmd: &mut dyn Reflect) {
let mut r = cmd.get_value().str();
ui.with_layout(
Layout::left_to_right(egui::Align::Min)
.with_main_wrap(true)
.with_main_justify(false),
|ui| {
ui.label(name);
ui.text_edit_singleline(&mut r);
cmd.set_value(ReflectValue::Str(r));
},
);
}
fn build_command_ui(ui: &mut Ui, name: &str, cmd: &mut dyn Reflect) {
// ui.separator();
match cmd.reflect_type() {
ReflectType::Value => match cmd.get_value() {
ReflectValue::None => {}
ReflectValue::U8(_) => build_command_ui_value_number(ui, name, cmd),
ReflectValue::U16(_) => build_command_ui_value_number(ui, name, cmd),
ReflectValue::U32(_) => build_command_ui_value_number(ui, name, cmd),
ReflectValue::U64(_) => build_command_ui_value_number(ui, name, cmd),
ReflectValue::I8(_) => build_command_ui_value_number(ui, name, cmd),
ReflectValue::I16(_) => build_command_ui_value_number(ui, name, cmd),
ReflectValue::I32(_) => build_command_ui_value_number(ui, name, cmd),
ReflectValue::I64(_) => build_command_ui_value_number(ui, name, cmd),
ReflectValue::Vec(_vec) => {
ui.with_layout(Layout::top_down(egui::Align::Min), |ui| {
ui.with_layout(
Layout::left_to_right(egui::Align::Min).with_main_wrap(true),
|ui| {
ui.label(name);
if ui.button("+").clicked() {
cmd.vec_add();
}
if ui.button("-").clicked() {
cmd.vec_remove();
}
},
);
if let Some(v) = cmd.as_vec() {
ui.with_layout(
Layout::left_to_right(egui::Align::Min).with_main_wrap(true),
|ui| {
for (i, cmd) in v.into_iter().enumerate() {
build_command_ui(ui, &format!("[{i}]"), cmd);
}
},
);
}
});
}
ReflectValue::Str(_) => {}
ReflectValue::Bool(_) => {}
},
ReflectType::Enumeration => {
// variants
ui.with_layout(
Layout::left_to_right(egui::Align::Min)
.with_main_wrap(true)
.with_main_justify(false),
|ui| {
for (i, (name, _v)) in cmd.variants().into_iter().enumerate() {
if ui.button(name).clicked() {
cmd.convert_variant(i);
}
}
},
);
// values
if let Some((name, cmd)) = cmd.unwrap_variant() {
build_command_ui(ui, name, cmd);
}
}
ReflectType::Structure => {
ui.with_layout(
Layout::left_to_right(egui::Align::Min)
.with_main_wrap(true)
.with_main_justify(false),
|ui| {
for (name, cmd) in cmd.fields().into_iter() {
build_command_ui(ui, name, cmd);
}
},
);
}
}
}
impl eframe::App for MyApp {
fn update(&mut self, ctx: &egui::Context, _frame: &mut eframe::Frame) {
egui::CentralPanel::default().show(ctx, |ui| {
ui.heading(format!("Command: {:?}", self.command));
if ui.button("X").clicked() {
self.command = Msg::None;
}
build_command_ui(ui, "Msg", &mut self.command);
// if self
// .command_path
// .as_bytes()
// .iter()
// .filter(|x| **x == b'.')
// .count()
// >= 2
// {
// for item in &self.tree.get_add_path(&self.command_path).children {
// for c in &item.children {
// if c.is_last() {
// if !self.values.contains_key(&c.name) {
// self.values.insert(c.name.clone(), String::new());
// }
// let sref = self.values.get_mut(&c.name).unwrap();
// ui.horizontal(|ui| {
// ui.label(format!("{} : {}", item.name, c.name.clone()));
// ui.text_edit_singleline(sref);
// });
// }
// }
// }
// if ui.button("send").clicked() {
// info!("sending {}", self.command_path);
// }
// } else {
// for item in &self.tree.get_add_path(&self.command_path).children {
// if ui.button(item.name.clone()).clicked() {
// self.command_path = format!("{}.{}", self.command_path, item.name);
// }
// }
// }
});
}
}
File diff suppressed because one or more lines are too long
+1
View File
@@ -8,6 +8,7 @@ embassy-usb-logger = { git = "https://github.com/embassy-rs/embassy.git" }
rfe = { path = "../../rfe", default-features = false, features = ["rp2040"] }
to = { path = "../../apps/to", default-features = false }
hs = { path = "../../apps/hs", default-features = false, features = ["rp2040"] }
example = { path = "../../apps/example" }
anyhow.workspace = true
anyhow.default_features = false
embedded-alloc = "0.6.0"
+6 -8
View File
@@ -15,10 +15,11 @@ mod app {
usb::{self, Driver},
};
use embedded_hal::digital::v2::OutputPin;
use example::Example;
use fugit::Duration;
use hs::{Hs, HsConfig, Rp2040Watchdog, StubSystemInfoGrabber};
use log::info;
use msg::{ExampleHk, Instance, MsgKind, MsgPacket, TargetMsg, TlmSetItem, ToTlmSet};
use msg::{Instance, MsgKind, MsgPacket, TargetMsg, TlmSetItem, ToTlmSet};
use rfe::{connector::Connector, Rate, *};
use rp_pico::hal::{
clocks,
@@ -69,19 +70,13 @@ mod app {
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 hk(&mut self, _rfe: &mut Rfe) {}
fn out_data(&mut self, _rfe: &mut Rfe) {}
@@ -198,7 +193,9 @@ mod app {
},
);
let mut to = To::new(&mut log_connector, tlmsets);
let mut example = Example::new();
let mut wd = Rp2040Watchdog::new(ctx.local.wd.take().unwrap());
let mut grabber = StubSystemInfoGrabber::new();
let mut hs = Hs::new(
HsConfig {
@@ -218,6 +215,7 @@ mod app {
instance.add_app("blink_app", &mut blink_app).unwrap();
instance.add_app("to", &mut to).unwrap();
instance.add_app("hs", &mut hs).unwrap();
instance.add_app("example", &mut example).unwrap();
let mut next_time = Mono::now() + Duration::<u64, 1, 1000000>::from_ticks(10000);
+8 -2
View File
@@ -5,9 +5,10 @@ edition = "2021"
[features]
default = []
std = ["dep:mio"]
std = ["dep:mio", "dep:mio-serial"]
rp2040 = ["dep:rp2040-hal", "dep:rp2040-pac"]
to_csv = []
reflect = ["dep:reflect"]
serde = ["dep:serde"]
[dependencies]
bincode.workspace = true
@@ -18,3 +19,8 @@ mio = { version = "1.0.2", features = ["net", "os-poll"], optional = true }
rp2040-hal = { workspace = true, optional = true }
rp2040-pac = { workspace = true, optional = true }
macros.path = "macros"
mio-serial = { workspace = true, optional = true }
serde = { version = "1.0.219", default-features = false, features = [
"derive",
], optional = true }
reflect = { workspace = true, optional = true }
+56
View File
@@ -0,0 +1,56 @@
# Real-time Framework for Embedded systems (RFE)
RFE is a framework for building real-time embedded applications in Rust. It provides a message-passing architecture for inter-application communication, time management, and scheduling at different rates.
## Features
- Message-passing architecture for inter-application communication
- Time management for both system and monotonic time
- Scheduling of applications at different rates (1Hz to 100Hz)
- Support for different platforms through feature flags
- Connectors for communication between instances (TCP, UDP, Memory)
## Usage
To use RFE, you need to create an RfeInstance, add applications to it, and then run it at 100Hz. Applications must implement the App trait.
```rust
use rfe::*;
struct MyApp;
impl App for MyApp {
fn init(&mut self, rfe: &mut Rfe) -> anyhow::Result<()> {
// Initialize the application
Ok(())
}
fn run(&mut self, rfe: &mut Rfe) {
// Run the application
}
fn hk(&mut self, rfe: &mut Rfe) {
// Generate housekeeping data
}
fn out_data(&mut self, rfe: &mut Rfe) {
// Generate output data
}
fn get_app_rate(&self) -> Rate {
Rate::Hz10 // Run at 10Hz
}
}
```
## Platform Support
RFE supports multiple platforms through feature flags:
- `std`: Standard library support (Unix, Windows)
- `rp2040`: Raspberry Pi Pico support
- `reflect`: Runtime type information for debugging
## License
This project is licensed under the MIT License - see the LICENSE file for details.
+40 -87
View File
@@ -31,6 +31,25 @@ pub fn kind_derive(input: TokenStream) -> TokenStream {
}
});
let kind_to_default_arms = variants.iter().map(|variant| {
let variant_name = &variant.ident;
match &variant.fields {
Fields::Unit => {
quote! {
#kind_name::#variant_name => #name::#variant_name,
}
}
Fields::Unnamed(_) => {
quote! {
#kind_name::#variant_name => #name::#variant_name(Default::default()),
}
}
Fields::Named(_) => {
panic!("Named fields not supported by Kind");
}
}
});
let kind_arms = variants.iter().map(|variant| {
let variant_name = &variant.ident;
match &variant.fields {
@@ -51,9 +70,11 @@ pub fn kind_derive(input: TokenStream) -> TokenStream {
});
let expanded = quote! {
#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[derive(Debug, Default, Copy, Clone, PartialEq, Eq, Hash, Encode, Decode)]
#[cfg_attr(feature = "reflect", derive(reflect::Reflect))]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub enum #kind_name {
#[default]
#(#kind_enum_arms)*
}
@@ -63,6 +84,12 @@ pub fn kind_derive(input: TokenStream) -> TokenStream {
#(#kind_enum_vec_arms)*
]
}
pub fn to_default(&self) -> #name {
match self {
#(#kind_to_default_arms)*
}
}
}
impl #name {
@@ -77,93 +104,19 @@ pub fn kind_derive(input: TokenStream) -> TokenStream {
TokenStream::from(expanded)
}
#[proc_macro_derive(ToCsv)]
pub fn to_csv_derive(input: TokenStream) -> TokenStream {
#[proc_macro_attribute]
pub fn msg(_attrs: TokenStream, input: TokenStream) -> TokenStream {
// Parse the input tokens into a syntax tree
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: alloc::vec::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) -> alloc::vec::Vec<String> {
extern crate alloc;
use alloc::format;
let mut values:alloc::vec::Vec<String> = alloc::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: alloc::vec::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) -> alloc::vec::Vec<String> {
extern crate alloc;
use alloc::format;
let mut values: alloc::vec::Vec<String> = alloc::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();
}
// Create a new version of the input with the additional derives
let output = quote! {
#[derive(Debug, Default, Clone, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "reflect", derive(reflect::Reflect))]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#input
};
TokenStream::from(expanded)
// Return the modified token stream
output.into()
}
+12
View File
@@ -0,0 +1,12 @@
[package]
name = "reflect"
version = "0.1.0"
edition = "2024"
[lib]
path = "src/reflect.rs"
name = "reflect"
[dependencies]
hashbrown.workspace = true
reflect_macros.path = "reflect_macros"
+14
View File
@@ -0,0 +1,14 @@
[package]
name = "reflect_macros"
version = "0.1.0"
edition = "2024"
[lib]
path = "src/reflect_macros.rs"
proc-macro = true
[dependencies]
proc-macro2 = "1.0.89"
quote = "1.0.37"
syn = { version = "2.0.87", features = ["full"] }
@@ -0,0 +1,196 @@
use proc_macro::TokenStream;
use quote::quote;
use syn::Fields;
use syn::{parse_macro_input, Data, DeriveInput};
#[proc_macro_derive(Reflect)]
pub fn reflect_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 field_arms = fields.iter().map(|field| {
let field_ident = field.clone().ident.unwrap();
let field_name = field_ident.to_string();
quote! { {
fields.push((#field_name, &mut self.#field_ident as &mut dyn reflect::Reflect));
} }
});
quote! {
impl reflect::Reflect for #name {
fn reflect_type(&self) -> reflect::ReflectType {
reflect::ReflectType::Structure
}
fn type_name(&self) -> &str {
#name_s
}
fn fields(&mut self) -> alloc::vec::Vec<(&str, &mut dyn reflect::Reflect)> {
let mut fields = alloc::vec::Vec::new();
#(#field_arms)*
fields
}
fn set_value(&mut self, value: reflect::ReflectValue) {
}
fn get_value(&self) -> reflect::ReflectValue {
reflect::ReflectValue::None
}
fn variants(&self) -> alloc::vec::Vec<(&'static str, alloc::boxed::Box<dyn reflect::Reflect>)> {
alloc::vec::Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn reflect::Reflect)> {
None
}
fn as_vec(&mut self) -> Option<alloc::vec::Vec<&mut dyn reflect::Reflect>> {
None
}
}
}
}
Data::Enum(data_enum) => {
let variants = data_enum.variants;
let convert_arms = variants.iter().enumerate().map(|(i, variant)| {
let variant_name = &variant.ident;
match &variant.fields {
Fields::Unit => {
quote! {
if _i == #i {
*self = Self::#variant_name;
}
}
}
Fields::Unnamed(_) => {
quote! {
if _i == #i {
*self = Self::#variant_name(Default::default());
}
}
}
Fields::Named(_) => {
panic!("Named fields not supported by Reflect");
}
}
});
let unwrap_arms = variants.iter().enumerate().map(|(_i, variant)| {
let variant_name = &variant.ident;
let variant_name_s = variant_name.to_string();
match &variant.fields {
Fields::Unit => {
quote! {
Self::#variant_name => None,
}
}
Fields::Unnamed(_) => {
quote! {
Self::#variant_name(v) => Some((#variant_name_s, v)),
}
}
Fields::Named(_) => {
panic!("Named fields not supported by Reflect");
}
}
});
let variant_arms = variants.iter().map(|variant| {
let variant_name = &variant.ident;
let variant_name_s = variant_name.to_string();
match &variant.fields {
Fields::Unit => {
quote! {
{
let name: &'static str = #variant_name_s;
let value: alloc::boxed::Box<dyn reflect::Reflect> = alloc::boxed::Box::new(Self::#variant_name);
variants.push((name, value));
}
}
}
Fields::Unnamed(_) => {
quote! {
{
let name: &'static str = #variant_name_s;
let value: alloc::boxed::Box<dyn reflect::Reflect> = alloc::boxed::Box::new(Self::#variant_name(Default::default()));
variants.push((name, value));
}
}
}
Fields::Named(_) => {
panic!("Named fields not supported by Reflect");
}
}
});
quote! {
impl reflect::Reflect for #name {
fn reflect_type(&self) -> reflect::ReflectType {
reflect::ReflectType::Enumeration
}
fn type_name(&self) -> &str {
#name_s
}
fn fields(&mut self) -> alloc::vec::Vec<(&str, &mut dyn reflect::Reflect)> {
alloc::vec::Vec::new()
}
fn set_value(&mut self, value: reflect::ReflectValue) {
}
fn get_value(&self) -> reflect::ReflectValue {
reflect::ReflectValue::None
}
fn variants(&self) -> alloc::vec::Vec<(&'static str, alloc::boxed::Box<dyn reflect::Reflect>)> {
let mut variants = alloc::vec::Vec::new();
#(#variant_arms)*
variants
}
fn convert_variant(&mut self, _i: usize) {
#(#convert_arms)*
}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn reflect::Reflect)> {
match self {
#(#unwrap_arms)*
}
}
fn as_vec(&mut self) -> Option<alloc::vec::Vec<&mut dyn reflect::Reflect>> {
None
}
}
}
}
Data::Union(_data_union) => {
return syn::Error::new_spanned(name, "Reflect not implemented for unions")
.to_compile_error()
.into();
}
};
TokenStream::from(expanded)
}
+650
View File
@@ -0,0 +1,650 @@
extern crate alloc;
use core::mem::zeroed;
use alloc::boxed::Box;
use alloc::string::String;
use alloc::vec::Vec;
use hashbrown::HashMap;
pub use reflect_macros::Reflect;
pub trait Reflect: std::fmt::Debug {
fn reflect_type(&self) -> ReflectType;
fn type_name(&self) -> &str;
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)>;
fn set_value(&mut self, value: ReflectValue);
fn get_value(&self) -> ReflectValue;
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)>;
fn convert_variant(&mut self, i: usize);
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)>;
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>>;
fn vec_add(&mut self) {}
fn vec_remove(&mut self) {}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReflectType {
Value,
Enumeration,
Structure,
}
pub enum ReflectValue {
None,
U8(u8),
U16(u16),
U32(u32),
U64(u64),
I8(i8),
I16(i16),
I32(i32),
I64(i64),
Vec(Vec<ReflectValue>),
Str(String),
Bool(bool),
}
impl ReflectValue {
pub fn signed(&self) -> i64 {
match self {
ReflectValue::None => 0,
ReflectValue::U8(v) => *v as i64,
ReflectValue::U16(v) => *v as i64,
ReflectValue::U32(v) => *v as i64,
ReflectValue::U64(v) => *v as i64,
ReflectValue::I8(v) => *v as i64,
ReflectValue::I16(v) => *v as i64,
ReflectValue::I32(v) => *v as i64,
ReflectValue::I64(v) => *v as i64,
ReflectValue::Vec(_) => 0,
ReflectValue::Str(s) => {
if let Ok(v) = s.parse() {
v
} else {
0
}
}
ReflectValue::Bool(b) => {
if *b {
1
} else {
0
}
}
}
}
pub fn unsigned(&self) -> u64 {
match self {
ReflectValue::None => 0,
ReflectValue::U8(v) => *v as u64,
ReflectValue::U16(v) => *v as u64,
ReflectValue::U32(v) => *v as u64,
ReflectValue::U64(v) => *v as u64,
ReflectValue::I8(v) => *v as u64,
ReflectValue::I16(v) => *v as u64,
ReflectValue::I32(v) => *v as u64,
ReflectValue::I64(v) => *v as u64,
ReflectValue::Vec(_) => 0,
ReflectValue::Str(s) => {
if let Ok(v) = s.parse() {
v
} else {
0
}
}
ReflectValue::Bool(b) => {
if *b {
1
} else {
0
}
}
}
}
pub fn str(self) -> String {
match self {
ReflectValue::Str(s) => s,
ReflectValue::None => "".to_string(),
ReflectValue::U8(v) => v.to_string(),
ReflectValue::U16(v) => v.to_string(),
ReflectValue::U32(v) => v.to_string(),
ReflectValue::U64(v) => v.to_string(),
ReflectValue::I8(v) => v.to_string(),
ReflectValue::I16(v) => v.to_string(),
ReflectValue::I32(v) => v.to_string(),
ReflectValue::I64(v) => v.to_string(),
ReflectValue::Vec(_v) => "".to_string(),
ReflectValue::Bool(v) => v.to_string(),
}
}
pub fn bool(&self) -> bool {
match self {
ReflectValue::None => false,
ReflectValue::U8(v) => *v == 0,
ReflectValue::U16(v) => *v == 0,
ReflectValue::U32(v) => *v == 0,
ReflectValue::U64(v) => *v == 0,
ReflectValue::I8(v) => *v == 0,
ReflectValue::I16(v) => *v == 0,
ReflectValue::I32(v) => *v == 0,
ReflectValue::I64(v) => *v == 0,
ReflectValue::Vec(_) => false,
ReflectValue::Str(_) => false,
ReflectValue::Bool(b) => *b,
}
}
pub fn vec(self) -> Vec<ReflectValue> {
if let ReflectValue::Vec(v) = self {
v
} else {
Vec::new()
}
}
}
pub fn path_set(reflect: &mut dyn Reflect, path: &str, value: ReflectValue) {
let mut next = reflect;
for p in path.split(".") {
let mut index = None;
for (i, f) in next.fields().iter().enumerate() {
if f.0 == p {
index = Some(i);
}
}
if let Some(i) = index {
next = next.fields().remove(i).1;
} else {
return;
}
}
next.set_value(value);
}
pub fn path_get<'a>(reflect: &'a mut dyn Reflect, path: &str) -> Option<&'a mut dyn Reflect> {
let mut next = reflect;
for p in path.split(".") {
let mut index = None;
for (i, f) in next.fields().iter().enumerate() {
if f.0 == p {
index = Some(i);
}
}
if let Some(i) = index {
next = next.fields().remove(i).1;
} else {
return None;
}
}
Some(next)
}
fn _flatten<'a>(reflect: &'a mut dyn Reflect) -> HashMap<String, &'a mut dyn Reflect> {
let mut list = HashMap::new();
match reflect.reflect_type() {
ReflectType::Value => match reflect.get_value() {
ReflectValue::Vec(_) => {
for (i, r) in reflect.as_vec().unwrap().into_iter().enumerate() {
let map = _flatten(r);
for (sn, v) in map {
list.insert(alloc::format!("[{i}].{}", sn), v);
}
}
}
_ => {
list.insert("".to_string(), reflect);
}
},
ReflectType::Enumeration => {
if let Some((name, r)) = reflect.unwrap_variant() {
let map = _flatten(r);
for (sn, v) in map {
list.insert(alloc::format!("{}.{}", name, sn), v);
}
}
}
ReflectType::Structure => {
for (name, r) in reflect.fields() {
let map = _flatten(r);
for (sn, v) in map {
list.insert(alloc::format!("{}.{}", name, sn), v);
}
}
}
}
list
}
pub fn flatten<'a>(reflect: &'a mut dyn Reflect) -> HashMap<String, &'a mut dyn Reflect> {
_flatten(reflect)
.into_iter()
.map(|(k, v)| (k[..k.len() - 1].to_string().replace(".[", "["), v))
.collect()
}
impl Reflect for u8 {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"u8"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.unsigned() as Self;
}
fn get_value(&self) -> ReflectValue {
ReflectValue::U8(*self)
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl Reflect for u16 {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"u16"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.unsigned() as Self;
}
fn get_value(&self) -> ReflectValue {
ReflectValue::U16(*self)
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl Reflect for u32 {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"u32"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.unsigned() as Self;
}
fn get_value(&self) -> ReflectValue {
ReflectValue::U32(*self)
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl Reflect for u64 {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"u64"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.unsigned() as Self;
}
fn get_value(&self) -> ReflectValue {
ReflectValue::U64(*self)
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl Reflect for i8 {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"i8"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.signed() as Self;
}
fn get_value(&self) -> ReflectValue {
ReflectValue::I8(*self)
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl Reflect for i16 {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"i16"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.signed() as Self;
}
fn get_value(&self) -> ReflectValue {
ReflectValue::I16(*self)
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl Reflect for i32 {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"i32"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.signed() as Self;
}
fn get_value(&self) -> ReflectValue {
ReflectValue::I32(*self)
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl Reflect for i64 {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"i64"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.signed() as Self;
}
fn get_value(&self) -> ReflectValue {
ReflectValue::I64(*self)
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl Reflect for String {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"str"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.str();
}
fn get_value(&self) -> ReflectValue {
ReflectValue::Str(self.to_string())
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl Reflect for bool {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"bool"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
*self = value.bool();
}
fn get_value(&self) -> ReflectValue {
ReflectValue::Bool(*self)
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec(&mut self) -> Option<Vec<&mut dyn Reflect>> {
None
}
}
impl<T: Reflect + Default> Reflect for Vec<T> {
fn reflect_type(&self) -> ReflectType {
ReflectType::Value
}
fn type_name(&self) -> &str {
"vec[]"
}
fn fields(&mut self) -> Vec<(&str, &mut dyn Reflect)> {
Vec::new()
}
fn set_value(&mut self, value: ReflectValue) {
if let ReflectValue::Vec(v) = value {
*self = v
.into_iter()
.map(|x| {
let mut t: T = unsafe { zeroed() };
t.set_value(x);
t
})
.collect();
}
}
fn get_value(&self) -> ReflectValue {
ReflectValue::Vec(self.iter().map(|x| x.get_value()).collect())
}
fn variants(&self) -> Vec<(&'static str, Box<dyn Reflect>)> {
Vec::new()
}
fn convert_variant(&mut self, _i: usize) {}
fn unwrap_variant(&mut self) -> Option<(&str, &mut dyn Reflect)> {
None
}
fn as_vec<'a>(&'a mut self) -> Option<Vec<&'a mut dyn Reflect>> {
Some(
self.iter_mut()
.map(|x| {
let v: &mut dyn Reflect = x;
v
})
.collect(),
)
}
fn vec_add(&mut self) {
self.push(Default::default());
}
fn vec_remove(&mut self) {
if self.len() > 0 {
self.remove(self.len() - 1);
}
}
}
+24
View File
@@ -1,11 +1,35 @@
/*!
* Connector module for RFE
*
* This module provides the Connector trait and implementations for various
* communication methods between RFE instances, including:
* - Memory connectors for inter-process communication
* - TCP connectors for network communication
* - UDP connectors for network communication
*/
use core::fmt::Debug;
use crate::msg::MsgPacket;
extern crate alloc;
use alloc::vec::Vec;
/// Connector trait for inter-instance communication
///
/// This trait defines the interface for connectors that enable communication
/// between RFE instances. Connectors are responsible for sending and receiving
/// messages between instances.
pub trait Connector: Debug {
/// Send messages to another instance
///
/// # Arguments
/// * `msgs` - The messages to send
fn send(&mut self, msgs: Vec<MsgPacket>);
/// Receive messages from another instance
///
/// # Returns
/// * `Some(Vec<MsgPacket>)` - If messages are available
/// * `None` - If no messages are available
fn recv(&mut self) -> Option<Vec<MsgPacket>>;
}
+79 -2
View File
@@ -1,23 +1,100 @@
#![no_std]
/*!
* Real-time Framework for Embedded systems (RFE)
*
* RFE is a framework for building real-time embedded applications in Rust.
* It provides a message-passing architecture for inter-application communication,
* time management, and scheduling at different rates.
*
* # Features
*
* - Message-passing architecture for inter-application communication
* - Time management for both system and monotonic time
* - Scheduling of applications at different rates (1Hz to 100Hz)
* - Support for different platforms through feature flags
* - Connectors for communication between instances (TCP, UDP, Memory)
*
* # Usage
*
* To use RFE, you need to create an RfeInstance, add applications to it,
* and then run it at 100Hz. Applications must implement the App trait.
*
* ```rust
* use rfe::*;
*
* struct MyApp;
*
* impl App for MyApp {
* fn init(&mut self, rfe: &mut Rfe) -> anyhow::Result<()> {
* // Initialize the application
* Ok(())
* }
*
* fn run(&mut self, rfe: &mut Rfe) {
* // Run the application
* }
*
* fn hk(&mut self, rfe: &mut Rfe) {
* // Generate housekeeping data
* }
*
* fn out_data(&mut self, rfe: &mut Rfe) {
* // Generate output data
* }
*
* fn get_app_rate(&self) -> Rate {
* Rate::Hz10 // Run at 10Hz
* }
* }
* ```
*/
#[cfg(feature = "std")]
extern crate std;
/// Connector module for inter-instance communication
pub mod connector;
use bincode::config::Configuration;
/// Message module for defining message types
pub mod msg;
/// Core RFE implementation
mod rfe;
pub use rfe::*;
#[cfg(feature = "to_csv")]
pub mod to_csv;
/// Reflection module for runtime type information
#[cfg(feature = "reflect")]
pub use reflect;
/// Macros for RFE
pub use macros;
/// Serial communication module
pub mod serial;
/// Time management module
pub mod time;
/// Utility functions and types
pub mod utils;
/// Bincode configuration for serialization
pub const BINCODE_CONFIG: Configuration = bincode::config::standard();
/// Macro for unwrapping a Result and printing an error message if it fails
///
/// # Arguments
/// * `$x` - The Result to unwrap
/// * `$msg` - The error message to print if the Result is an Err
///
/// # Example
/// ```rust
/// use rfe::unwrap_print_err;
/// use log::error;
///
/// fn main() {
/// let result: Result<(), &str> = Err("error");
/// unwrap_print_err!(result, "Failed to do something");
/// }
/// ```
#[macro_export]
macro_rules! unwrap_print_err {
($x:expr, $msg: tt) => {
+9 -12
View File
@@ -1,17 +1,14 @@
#[cfg(feature = "to_csv")]
use crate::to_csv::ToCsv;
use crate::utils::PerfData;
use bincode::{Decode, Encode};
#[cfg(feature = "to_csv")]
use macros::ToCsv;
extern crate alloc;
use alloc::string::String;
use alloc::vec::Vec;
use macros::msg;
use super::{TlmSetId, TlmSetItem};
#[derive(Debug, Clone, PartialEq, Eq, Hash, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Eq, Hash)]
pub struct DsTlmSet {
pub items: Vec<TlmSetItem>,
pub id: TlmSetId,
@@ -19,24 +16,24 @@ pub struct DsTlmSet {
pub path: String,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq)]
pub struct DsHk {
pub perf: PerfData,
pub counter: u32,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq)]
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))]
#[msg]
pub enum DsCmd {
#[default]
Noop,
Reset,
CloseAll,
+8 -10
View File
@@ -1,27 +1,25 @@
extern crate alloc;
#[cfg(feature = "to_csv")]
use crate::to_csv::ToCsv;
use crate::utils::PerfData;
use bincode::{Decode, Encode};
#[cfg(feature = "to_csv")]
use macros::ToCsv;
use macros::msg;
#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq)]
pub struct ExampleHk {
pub perf: PerfData,
pub counter: u32,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq)]
pub struct ExampleOutData {
pub counter: u32,
}
#[derive(Debug, Clone, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq)]
pub enum ExampleCmd {
#[default]
Noop,
Reset,
}
+8 -10
View File
@@ -1,14 +1,11 @@
#[cfg(feature = "to_csv")]
use crate::to_csv::ToCsv;
use crate::utils::PerfData;
use bincode::{Decode, Encode};
#[cfg(feature = "to_csv")]
use macros::ToCsv;
use macros::msg;
extern crate alloc;
use alloc::vec::Vec;
#[derive(Debug, Default, Clone, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
pub struct HsHk {
pub perf: PerfData,
pub counter: u32,
@@ -22,15 +19,16 @@ pub struct HsHk {
pub fs_usage_enabled: bool,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq)]
pub struct HsOutData {
pub counter: u32,
}
#[derive(Debug, Clone, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq)]
pub enum HsCmd {
#[default]
Noop,
Reset,
WatchdogEnableManual(bool),
+16 -19
View File
@@ -3,7 +3,7 @@ extern crate alloc;
use alloc::string::String;
use alloc::vec::Vec;
use bincode::{Decode, Encode};
use macros::Kind;
use macros::{msg, Kind};
mod example;
pub use example::*;
@@ -15,13 +15,9 @@ mod ds;
pub use ds::*;
use crate::time::Timestamp;
#[cfg(feature = "to_csv")]
use crate::to_csv::ToCsv;
#[cfg(feature = "to_csv")]
use macros::ToCsv;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq, Hash)]
pub struct TargetMsg {
// instance is FROM for tlm, is TO for cmds
pub instance: Instance,
@@ -34,18 +30,21 @@ impl TargetMsg {
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq, Hash)]
pub enum Instance {
#[default]
None,
All,
Other,
Example,
Example2,
}
#[derive(Debug, Clone, PartialEq, Kind, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Kind)]
pub enum Msg {
#[default]
None,
SubRequest,
SubList(SubList),
@@ -65,8 +64,7 @@ pub enum Msg {
ToCmd(ToCmd),
}
#[derive(Debug, Clone, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
pub struct MsgPacket {
pub instance: Instance,
pub msg: Msg,
@@ -91,24 +89,23 @@ impl MsgPacket {
}
}
#[derive(Debug, Clone, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
pub struct SubList {
pub subs: Vec<TargetMsg>,
}
#[derive(Debug, Clone, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
pub struct ReinitAppCmd {
app_name: String,
}
pub type TlmSetId = u16;
#[derive(Debug, Clone, PartialEq, Eq, Hash, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq, Hash)]
pub struct TlmSetItem {
pub target: TargetMsg,
/// decimation 0 means sends every msg, decimation 1 means sends every other, etc
pub decimation: u16,
pub counter: u16,
}
+10 -12
View File
@@ -1,38 +1,36 @@
#[cfg(feature = "to_csv")]
use crate::to_csv::ToCsv;
use crate::utils::PerfData;
use bincode::{Decode, Encode};
#[cfg(feature = "to_csv")]
use macros::ToCsv;
use macros::msg;
extern crate alloc;
use alloc::vec::Vec;
use super::{TlmSetId, TlmSetItem};
#[derive(Debug, Clone, PartialEq, Eq, Hash, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Eq, Hash)]
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))]
#[msg]
#[derive(Copy, Eq)]
pub struct ToHk {
pub perf: PerfData,
pub counter: u32,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq)]
pub struct ToOutData {
pub counter: u32,
}
#[derive(Debug, Clone, PartialEq, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Eq)]
pub enum ToCmd {
#[default]
Noop,
Reset,
AddTlmSet(ToTlmSet),
View File
+252 -34
View File
@@ -1,3 +1,12 @@
/*!
* Real-time Framework for Embedded systems (RFE)
*
* This module provides the core functionality for the RFE framework:
* - Application scheduling at different rates (1Hz to 100Hz)
* - Message passing between applications
* - Time management for both system and monotonic time
* - Connector management for inter-instance communication
*/
extern crate alloc;
use core::cell::RefCell;
@@ -13,99 +22,207 @@ use crate::{
time::{TimeData, TimeDriver},
};
/// Housekeeping data trait
///
/// This trait is used to mark types that can be used as housekeeping data.
/// Housekeeping data is used to monitor the health and status of applications.
pub trait Hk: Sized + Clone + Copy + 'static + Send + Sync {}
/// Blanket implementation for all types that meet the requirements
impl<T> Hk for T where T: Sized + Clone + Copy + 'static + Send + Sync {}
/// Output data trait
///
/// This trait is used to mark types that can be used as output data.
/// Output data is the primary data produced by applications.
pub trait OutData: Sized + Clone + Copy + 'static + Send + Sync {}
/// Blanket implementation for all types that meet the requirements
impl<T> OutData for T where T: Sized + Clone + Copy + 'static + Send + Sync {}
/// Application execution rates
///
/// This enum defines the rates at which applications can be scheduled.
/// The RFE framework runs at 100Hz, and applications can be scheduled
/// at various rates derived from this base rate.
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum Rate {
/// 1Hz (once per second)
Hz1,
/// 5Hz (5 times per second)
Hz5,
/// 10Hz (10 times per second)
Hz10,
/// 20Hz (20 times per second)
Hz20,
/// 50Hz (50 times per second)
Hz50,
/// 100Hz (100 times per second, every tick)
Hz100,
}
/// Reference to an RfeTime instance
///
/// This type alias is used to share a single time reference between multiple RFE instances.
type RfeTimeRef<'a> = Rc<RefCell<RfeTime<'a>>>;
/// Time management for RFE
///
/// This struct manages time for the RFE framework, providing both system time
/// and monotonic time through a TimeDriver implementation.
pub struct RfeTime<'a> {
/// Time data including scheduler counter and time offset
time_data: TimeData,
/// Driver for time operations
time_driver: &'a dyn TimeDriver,
}
/// Main RFE instance
///
/// This struct represents a single RFE instance, which can contain multiple applications
/// and connectors. It manages the scheduling of applications and the routing of messages.
pub struct RfeInstance<'a> {
/// List of applications registered with this instance
app_list: HashMap<&'a str, AppRef<'a>>,
/// Reference to the time management
time: RfeTimeRef<'a>,
/// The instance identifier
#[allow(dead_code)]
instance: Instance,
/// List of connectors for inter-instance communication
connectors: Vec<ConnectorState<'a>>,
/// Scheduler counter, incremented on each tick
sch_counter: u64,
}
/// Reference to an application
///
/// This struct holds a reference to an application and its associated RFE instance.
/// It also stores the rates at which the application should be run.
pub struct AppRef<'a> {
/// Reference to the application
app: &'a mut dyn App,
/// Rate at which the application's run method should be called
app_rate: Rate,
/// Rate at which the application's out_data method should be called
out_data_rate: Rate,
/// Rate at which the application's hk method should be called
hk_rate: Rate,
/// RFE instance for this application
rfe: Rfe<'a>,
}
/// Core RFE interface for applications
///
/// This struct provides the interface for applications to interact with the RFE framework.
/// It handles message subscription, sending, and receiving, as well as time management.
pub struct Rfe<'a> {
/// Messages this application is subscribed to
subscriptions: HashSet<TargetMsg>,
/// Messages to be sent to other applications or connectors
msgs_to_send: Vec<MsgPacket>,
msgs_recevied: VecDeque<MsgPacket>,
/// Messages received from other applications or connectors
msgs_received: VecDeque<MsgPacket>,
/// The instance this RFE belongs to
instance: Instance,
/// Reference to the time management
time: RfeTimeRef<'a>,
/// Flag indicating if subscriptions have been updated
subs_updated: bool,
}
/// State of a connector
///
/// This struct holds a reference to a connector and its associated state.
/// It tracks subscriptions and subscription request timing.
#[derive(Debug)]
pub struct ConnectorState<'a> {
/// Reference to the connector
connector: &'a mut dyn Connector,
/// Messages this connector is subscribed to
subscriptions: HashSet<TargetMsg>,
/// Flag indicating if subscriptions have been received
subs_received: bool,
/// Scheduler counter value when subscriptions were last requested
subs_last_requested: u64,
}
impl<'a> Rfe<'a> {
/// Creates a new RFE instance
///
/// # Arguments
/// * `instance` - The instance this RFE belongs to
/// * `time` - Reference to the time management
pub fn new(instance: Instance, time: RfeTimeRef<'a>) -> Self {
Self {
subscriptions: HashSet::new(),
msgs_to_send: Vec::new(),
msgs_recevied: VecDeque::new(),
msgs_received: VecDeque::new(),
instance,
subs_updated: false,
time,
}
}
/// Returns the instance this RFE belongs to
pub fn get_instance(&self) -> Instance {
return self.instance;
self.instance
}
/// Subscribe to a message
///
/// This method subscribes the application to a specific message type.
/// The application will receive messages of this type from other applications
/// and connectors.
///
/// # Arguments
/// * `msg` - The message type to subscribe to
pub fn subscribe(&mut self, msg: TargetMsg) {
self.subscriptions.insert(msg);
self.subs_updated = true;
}
/// Subscribe to multiple messages
///
/// This method subscribes the application to multiple message types.
/// The application will receive messages of these types from other applications
/// and connectors.
///
/// # Arguments
/// * `msgs` - The message types to subscribe to
pub fn subscribe_all<T: IntoIterator<Item = TargetMsg>>(&mut self, msgs: T) {
self.subscriptions.extend(msgs.into_iter());
self.subs_updated = true;
}
/// Unsubscribe from a message
///
/// This method unsubscribes the application from a specific message type.
/// The application will no longer receive messages of this type.
///
/// # Arguments
/// * `msg` - The message type to unsubscribe from
pub fn unsubscribe(&mut self, msg: &TargetMsg) {
self.subscriptions.remove(msg);
self.subs_updated = true;
}
/// Unsubscribe from all messages
///
/// This method unsubscribes the application from all message types.
/// The application will no longer receive any messages.
pub fn unsubscribe_all(&mut self) {
self.subscriptions.clear();
self.subs_updated = true;
}
/// Send a message
///
/// This method sends a message from this application to other applications
/// and connectors that are subscribed to this message type.
///
/// # Arguments
/// * `msg` - The message to send
pub fn send(&mut self, msg: Msg) {
self.msgs_to_send.push(MsgPacket::new(
self.get_instance(),
@@ -114,41 +231,110 @@ impl<'a> Rfe<'a> {
));
}
/// Send a command to a specific instance
///
/// This method sends a command to a specific instance.
/// The command will be received by applications in that instance
/// that are subscribed to this message type.
///
/// # Arguments
/// * `msg` - The command to send
/// * `target` - The target instance to send the command to
pub fn send_cmd(&mut self, msg: Msg, target: Instance) {
self.msgs_to_send
.push(MsgPacket::new(target, msg, self.get_system_time()));
}
/// Posts a message to this RFE's receive queue
///
/// # Arguments
/// * `msg` - The message to post
pub fn post_message(&mut self, msg: MsgPacket) {
self.msgs_recevied.push_back(msg);
self.msgs_received.push_back(msg);
}
/// Receives a message from this RFE's receive queue
///
/// # Returns
/// * `Some(MsgPacket)` - If a message is available
/// * `None` - If no message is available
pub fn recv(&mut self) -> Option<MsgPacket> {
self.msgs_recevied.pop_front()
self.msgs_received.pop_front()
}
/// time starting from power on or program start
/// Get mission elapsed time
///
/// This method returns the time in microseconds since power on or program start.
/// It is useful for measuring durations and scheduling events.
///
/// # Returns
/// * Time in microseconds since power on or program start
pub fn get_met_time(&self) -> u64 {
let time = self.time.borrow();
time.time_driver.get_monotonic_time(time.time_data)
}
/// Time in microseconds relative to system epoch
/// Get system time
///
/// This method returns the time in microseconds relative to the system epoch.
/// It is useful for timestamping events and correlating with external systems.
///
/// # Returns
/// * Time in microseconds relative to the system epoch
pub fn get_system_time(&self) -> u64 {
let time = self.time.borrow();
time.time_driver.get_system_time(time.time_data)
}
}
/// Application trait for RFE applications
///
/// This trait defines the interface for applications to interact with the RFE framework.
/// Applications must implement this trait to be scheduled by the RFE framework.
pub trait App {
/// Initialize the application
///
/// This method is called once when the application is added to the RFE instance.
/// Use this method to set up subscriptions and initialize the application state.
fn init(&mut self, rfe: &mut Rfe) -> Result<()>;
/// Run the application
///
/// This method is called at the rate specified by `get_app_rate()`.
/// Use this method to perform the main application logic.
fn run(&mut self, rfe: &mut Rfe);
/// Generate housekeeping data
///
/// This method is called at the rate specified by the RFE instance.
/// Use this method to generate housekeeping data for telemetry.
fn hk(&mut self, rfe: &mut Rfe);
/// Generate output data
///
/// This method is called at the rate specified by the RFE instance.
/// Use this method to generate output data for telemetry.
fn out_data(&mut self, rfe: &mut Rfe);
/// Get the application rate
///
/// This method returns the rate at which the application should be run.
fn get_app_rate(&self) -> Rate;
}
impl<'a> RfeInstance<'a> {
/// Create a new RFE instance
///
/// This method creates a new RFE instance with the specified instance identifier
/// and time driver. The instance identifier is used to identify this instance
/// when communicating with other instances.
///
/// # Arguments
/// * `instance` - The instance identifier
/// * `time_driver` - The time driver to use for time management
///
/// # Returns
/// * A new RFE instance
pub fn new(instance: Instance, time_driver: &'a dyn TimeDriver) -> Self {
let time = Rc::new(RefCell::new(RfeTime {
time_data: TimeData {
@@ -166,32 +352,50 @@ impl<'a> RfeInstance<'a> {
}
}
/// Add an application to the RFE instance
///
/// # Arguments
/// * `name` - The name of the application
/// * `app` - The application to add
///
/// # Returns
/// * `Ok(())` - If the application was added successfully
/// * `Err(...)` - If the application could not be added
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"
"Failed to add app '{name}': an app with that name already exists"
));
}
let app_rate = app.get_app_rate();
self.app_list.insert(
name,
AppRef {
app: app,
app_rate: app_rate,
app,
app_rate,
hk_rate: Rate::Hz1,
out_data_rate: app_rate,
rfe: Rfe::new(self.instance, self.time.clone()),
},
);
// Initialize the application
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}");
error!("App '{name}' failed to initialize: {e}");
}
return Ok(());
Ok(())
}
/// Add a connector to the RFE instance
///
/// Connectors are used to communicate with other RFE instances.
/// They can be used to send and receive messages between instances.
///
/// # Arguments
/// * `connector` - The connector to add
pub fn add_connector(&mut self, connector: &'a mut dyn Connector) {
self.connectors.push(ConnectorState {
connector,
@@ -201,40 +405,31 @@ impl<'a> RfeInstance<'a> {
});
}
/// Expected to be called at 100Hz
/// Run the RFE instance
///
/// This method should be called at 100Hz to schedule applications and process messages.
/// It runs applications, collects messages, and routes them to the appropriate destinations.
pub fn run(&mut self) {
let mut msgs = Vec::new();
// Run applications and collect messages
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)
{
// Check if the application should run at this tick
if Self::should_run_at_rate(self.sch_counter, app.app_rate) {
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)
{
// Check if housekeeping should run at this tick
if Self::should_run_at_rate(self.sch_counter, app.hk_rate) {
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)
{
// Check if output data should run at this tick
if Self::should_run_at_rate(self.sch_counter, app.out_data_rate) {
app.app.out_data(&mut app.rfe);
}
// Collect messages from the application
let new_msgs = core::mem::take(&mut app.rfe.msgs_to_send);
msgs.extend(new_msgs);
}
@@ -365,7 +560,30 @@ impl<'a> RfeInstance<'a> {
self.sch_counter += 1;
}
/// Helper method to determine if a task should run at the given rate
///
/// # Arguments
/// * `rate` - The rate to check
///
/// # Returns
/// * `true` - If the task should run at this tick
/// * `false` - If the task should not run at this tick
fn should_run_at_rate(sch_counter: u64, rate: Rate) -> bool {
match rate {
Rate::Hz100 => true,
Rate::Hz50 => sch_counter % 2 == 0,
Rate::Hz20 => sch_counter % 5 == 0,
Rate::Hz10 => sch_counter % 10 == 0,
Rate::Hz5 => sch_counter % 20 == 0,
Rate::Hz1 => sch_counter % 100 == 0,
}
}
#[cfg(feature = "std")]
/// Start the RFE instance
///
/// This method starts the RFE instance and runs it at 100Hz.
/// It blocks the current thread and never returns.
pub fn start(&mut self) {
use core::time::Duration;
use std::{thread::sleep, time::Instant};
+55
View File
@@ -0,0 +1,55 @@
use anyhow::Result;
pub trait SerialPort {
fn open(&mut self) -> Result<()>;
fn write(&mut self, bytes: &[u8]) -> Result<usize>;
fn read(&mut self, bytes: &mut [u8]) -> Result<usize>;
}
#[cfg(feature = "std")]
mod serial_std {
use std::io::{Read, Write};
use anyhow::anyhow;
use mio_serial::{SerialPortBuilder, SerialPortBuilderExt, SerialStream};
pub struct StdSerialPort {
pub config: SerialPortBuilder,
serial: Option<SerialStream>,
}
impl StdSerialPort {
pub fn new(config: SerialPortBuilder) -> Self {
Self {
config,
serial: None,
}
}
}
impl super::SerialPort for StdSerialPort {
fn open(&mut self) -> anyhow::Result<()> {
self.serial = None;
self.serial = Some(self.config.clone().open_native_async()?);
return Ok(());
}
fn write(&mut self, bytes: &[u8]) -> anyhow::Result<usize> {
Ok(self
.serial
.as_mut()
.ok_or(anyhow!("port not open"))?
.write(bytes)?)
}
fn read(&mut self, bytes: &mut [u8]) -> anyhow::Result<usize> {
Ok(self
.serial
.as_mut()
.ok_or(anyhow!("port not open"))?
.read(bytes)?)
}
}
}
#[cfg(feature = "std")]
pub use serial_std::*;
+44 -2
View File
@@ -1,17 +1,59 @@
/*!
* Time management module for RFE
*
* This module provides time management functionality for the RFE framework,
* including:
* - Timestamp type for representing time
* - TimeData struct for storing time-related data
* - TimeDriver trait for platform-specific time implementations
* - Various TimeDriver implementations for different platforms
*/
/// Microseconds timestamp
///
/// This type represents time in microseconds, either as a duration or
/// as an absolute time relative to some epoch.
pub type Timestamp = u64;
/// Time data for RFE
///
/// This struct stores time-related data for the RFE framework, including
/// the scheduler counter and time offset.
#[derive(Debug, Clone, Copy)]
pub struct TimeData {
/// Scheduler counter, incremented on each tick
pub sch_counter: u64,
/// Time offset in microseconds
pub time_offset: Timestamp,
}
/// Time driver trait
///
/// This trait defines the interface for platform-specific time implementations.
/// It provides methods for getting both system time and monotonic time.
pub trait TimeDriver {
/// Time in microseconds relative to system epoch
/// Get system time
///
/// This method returns the time in microseconds relative to the system epoch.
/// It is useful for timestamping events and correlating with external systems.
///
/// # Arguments
/// * `time_data` - The time data to use
///
/// # Returns
/// * Time in microseconds relative to the system epoch
fn get_system_time(&self, time_data: TimeData) -> Timestamp;
/// Time in microseconds since program start or power on
/// Get monotonic time
///
/// This method returns the time in microseconds since program start or power on.
/// It is useful for measuring durations and scheduling events.
///
/// # Arguments
/// * `time_data` - The time data to use
///
/// # Returns
/// * Time in microseconds since program start or power on
fn get_monotonic_time(&self, time_data: TimeData) -> Timestamp;
}
-151
View File
@@ -1,151 +0,0 @@
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
}
+4 -7
View File
@@ -1,15 +1,12 @@
use crate::time::Timestamp;
use crate::Rfe;
use crate::time::Timestamp;
use bincode::{Decode, Encode};
use macros::msg;
extern crate alloc;
#[cfg(feature = "to_csv")]
use crate::to_csv::ToCsv;
#[cfg(feature = "to_csv")]
use macros::ToCsv;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash, Encode, Decode)]
#[cfg_attr(feature = "to_csv", derive(ToCsv))]
#[msg]
#[derive(Copy, Eq, Hash)]
pub struct PerfData {
enter_time: Timestamp,
elapsed: u32,
+3 -1
View File
@@ -11,5 +11,7 @@ path = "src/decom.rs"
anyhow.workspace = true
log.workspace = true
simple_logger.workspace = true
rfe = { path = "../../rfe", features = ["to_csv"] }
rfe = { path = "../../rfe", features = ["serde"] }
bincode = { workspace = true, features = ["std"] }
serde_json = "1.0.140"
clap = { version = "4.5.35", features = ["derive"] }
+20 -30
View File
@@ -1,7 +1,6 @@
#![feature(bufreader_peek)]
use std::{
collections::HashMap,
env::args,
fs::{read_dir, OpenOptions},
io::{BufReader, Write},
path::PathBuf,
@@ -12,10 +11,20 @@ use std::{
use anyhow::Result;
use bincode::decode_from_std_read;
use log::*;
use rfe::ToCsvClean;
use rfe::{msg::MsgPacket, BINCODE_CONFIG};
use simple_logger::SimpleLogger;
use clap::Parser;
// Decom DS recorded data files
#[derive(Parser, Debug)]
struct CliArgs {
/// The path to the file or folder containing the DS data
path: String,
/// The path to the folder to output the decommed data
out_dir: String,
}
fn decom_file(file_path: String, out_dir: String) -> Result<()> {
info!("decomming {file_path}");
let f = OpenOptions::new().read(true).open(file_path)?;
@@ -30,40 +39,25 @@ fn decom_file(file_path: String, out_dir: String) -> Result<()> {
};
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()
.join(format!("{:?}.json", msg.msg.kind()).to_lowercase());
let 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())?;
let values = serde_json::to_string(&msg)?;
w.write(values.as_bytes())?;
w.write("\n".as_bytes())?;
}
}
}
return Ok(());
}
@@ -100,13 +94,9 @@ fn decom_task(folder: String, out_dir: String) -> Result<()> {
fn main() -> Result<()> {
SimpleLogger::new().init().unwrap();
let args = CliArgs::parse();
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)?;
decom_task(args.path, args.out_dir)?;
info!("finished decom");
return Ok(());