feat(yrobot-panel): 实现电机控制上位机

完成 Rust/egui 上位机主流程:BLE/USB 传输、协议编解码、控制/实时/统计/CSV、UI 与验证脚本。\n\n保留真实设备联调和长稳实测清单,避免把未验证的硬件结果当成交付完成。
This commit is contained in:
2026-06-12 17:46:29 +08:00
commit dd442f0f22
32 changed files with 8527 additions and 0 deletions
+9
View File
@@ -0,0 +1,9 @@
.agents/
.claude/
.codegraph/
x86_64-w64-mingw32-dlltool
# Added by cargo
/target
Generated
+4871
View File
File diff suppressed because it is too large Load Diff
+21
View File
@@ -0,0 +1,21 @@
[package]
name = "yrobot-motor-control-panel"
version = "0.1.0"
edition = "2021"
[dependencies]
eframe = { version = "0.31", features = ["default_fonts", "glow"] }
egui_plot = "0.31"
log = "0.4"
env_logger = "0.11"
tokio = { version = "1", features = ["rt", "macros", "sync", "time"] }
tokio-util = "0.7"
btleplug = "0.11"
serialport = "4"
crc = "3"
futures = "0.3"
serde = { version = "1", features = ["derive"] }
serde_json = "1"
csv = "1"
thiserror = "2"
uuid = "1"
+78
View File
@@ -0,0 +1,78 @@
# YRobot Motor Control Panel
Rust/egui desktop tool for YRobot motor debugging. It supports direct BLE and USB Dongle transport, motor control commands, realtime debug plots, CSV recording, and session statistics.
## Build
Windows executable:
```bash
cargo build --release --target x86_64-pc-windows-gnu
```
Output:
```text
target/x86_64-pc-windows-gnu/release/yrobot-motor-control-panel.exe
```
Run tests for the Windows target:
```bash
cargo test --target x86_64-pc-windows-gnu
```
Run the full local verification gate:
```bash
./scripts/verify.sh
```
## Minimal Project Files
To hand the project to another developer, keep these tracked files:
- `Cargo.toml`
- `Cargo.lock`
- `docs/`
- `README.md`
- `scripts/`
- `src/`
Do not include local/generated directories such as `target/`, `.agents/`, `.claude/`, or `.codegraph/`.
## Use
1. Select `BLE` or `USB Dongle`.
2. Select `Left`, `Right`, or `Neutral` before connecting.
3. BLE: click `Scan / Refresh`, select a device, then `Connect BLE`. Use `Stop Scan` to stop discovery.
4. USB Dongle: refresh/select a serial port, then `Open USB`. The serial session uses `115200 8N1`.
5. After connection, use motor controls, realtime debug, CSV recording, and status/statistics panels.
Heartbeat starts only after BLE notify subscription or USB serial open succeeds, and stops when the session is disconnected.
CSV recording is flushed and closed when recording is stopped or the session ends.
The status panel shows live metrics and the latest error summary for the current session.
Closing the app cancels scans, disconnects the active session, and closes any active CSV recorder.
## Protocol Notes
- BLE uses Nordic UART characteristics:
- TX: `6e400002-b5a3-f393-e0a9-e50e24dcca9e`
- RX notify: `6e400003-b5a3-f393-e0a9-e50e24dcca9e`
- USB Dongle downlink wraps motor payload as `5A A5 + DEV(0x31) + LEN + PAYLOAD + CRC8/MAXIM`.
- USB Dongle uplink extracts `seq_id(2B LE)` and `dt_ms(2B LE)` before motor protocol parsing.
## Verification
Current automated checks cover data model defaults/serde, protocol encode/decode, heartbeat cancel/error behavior, parameter validation, tuning packet reassembly, debug-field error marking, communication-event logging, CSV writing/no-op/error paths, session-end/app-close recorder cleanup, host interval statistics, USB `seq_id`/`dt_ms` statistics, and simulated BLE/USB receive paths into chart/statistics.
Hardware validation still requires a real YRobot motor and USB Dongle:
- BLE scan/connect/notify/write.
- 2s heartbeat under real connection.
- Motor parameter set/query.
- Realtime debug stream at device rate.
- USB Dongle serial open, wrapping/unwrapping, `seq_id`/`dt_ms` diagnostics.
- Long-run stability on the target machine.
Use [docs/hardware-validation.md](docs/hardware-validation.md) to record final hardware acceptance.
+128
View File
@@ -0,0 +1,128 @@
# Hardware Validation Checklist
Use this checklist with a real YRobot motor and USB Dongle before marking the project fully complete.
## Environment
- Windows 10+ machine.
- Built executable: `target/x86_64-pc-windows-gnu/release/yrobot-motor-control-panel.exe`.
- One YRobot BLE motor.
- One USB Dongle and matching serial driver, if USB mode is being validated.
Record:
- App version/build time:
- Windows version:
- Motor identifier:
- USB Dongle port:
- Tester:
- Date/time:
## BLE Direct Mode
1. Start the app and select `BLE`.
2. Select `Left`, `Right`, or `Neutral`.
3. Click `Scan / Refresh`.
4. Confirm target device appears with name/ref/RSSI and preferred marker when applicable.
5. Click `Stop Scan` and confirm discovery stops.
6. Select the device and click `Connect BLE`.
7. Confirm status becomes `Connected` only after BLE notify is ready.
8. Keep the session open for at least 30 seconds and confirm no heartbeat error appears.
9. Send:
- Query Current
- Query Mode
- Set Params
- Enable Realtime
- Disable Realtime
10. Confirm TX events and RX events appear in the event log.
11. Confirm realtime samples appear in chart/status when debug stream is enabled.
12. Click `Disconnect` and confirm connection closes and heartbeat stops.
Pass criteria:
- No app crash.
- No stale scanning after stop/disconnect.
- Commands are logged.
- Realtime data parses without blocking valid fields.
- Latest error is empty or explains any hardware/protocol issue.
## USB Dongle Mode
1. Start the app and select `USB Dongle`.
2. Refresh ports and select the Dongle serial port.
3. Click `Open USB`.
4. Confirm status becomes `Connected` with `115200 8N1`.
5. Keep the session open for at least 30 seconds and confirm heartbeat starts only after serial open.
6. Send:
- Query Current
- Query Mode
- Set Params
- Enable Realtime
- Disable Realtime
7. Confirm USB RX events include `DEV`, `LEN`, `seq_id`, and `dt_ms`.
8. Confirm USB statistics show:
- seq gaps
- out-of-order count
- USB dt avg/max
- host interval avg/max
9. Disconnect and confirm serial resources are released.
Pass criteria:
- No app crash.
- Serial open failure is shown as a recoverable error.
- Downlink uses `DEV=0x31`.
- Uplink strips `seq_id`/`dt_ms` before motor parser.
- USB metadata is visible in event log/status/CSV.
## CSV Recording
1. Connect through BLE or USB.
2. Set a writable CSV path.
3. Start recording.
4. Enable realtime debug and collect at least 100 samples.
5. Stop recording.
6. Open the CSV and confirm it contains:
- timestamp
- connection mode
- target id
- side
- raw debug string
- parsed current/velocity/position/voltage/temperature
- USB `seq_id`/`dt_ms` for USB mode
7. Repeat with recording active, then disconnect. Confirm CSV is flushed and readable.
Pass criteria:
- CSV is not created before recording starts.
- CSV write errors are shown and recording stops.
- Session end closes the file cleanly.
## Stability
Run one session for at least 1 hour before release, and 8 hours for final acceptance.
Record every 15 minutes:
- sample count
- app memory usage
- latest error
- drop rate
- USB seq gaps/out-of-order, if USB mode
Pass criteria:
- UI remains responsive.
- Control, display, CSV, and statistics work concurrently.
- No obvious memory growth above 100 MB over the acceptance run.
- Disconnect closes scan/session/CSV resources.
## Result
- BLE Direct: Pass / Fail
- USB Dongle: Pass / Fail
- CSV Recording: Pass / Fail
- Stability: Pass / Fail
Notes:
+20
View File
@@ -0,0 +1,20 @@
#!/usr/bin/env bash
set -euo pipefail
TARGET="${1:-x86_64-pc-windows-gnu}"
cargo test --target "$TARGET"
cargo build --release --target "$TARGET"
if rg -n -i 'not implemented|not wired|BLE scan not wired|Serial not implemented|BLE not implemented|TODO|FIXME|placeholder|dummy|stub' \
-g 'src/**' \
-g 'Cargo.toml' \
-g 'README.md' \
-g 'docs/**' \
-g '.agents/specs/yrobot-motor-control-panel/**' \
.; then
echo "Found unfinished marker text" >&2
exit 1
fi
echo "Verification passed for target: $TARGET"
+522
View File
@@ -0,0 +1,522 @@
use crate::models::{ConnectionMode, DeviceCandidate, SessionState, SideKey, UiSessionViewState};
use crate::protocol::motor_parser::MotorProtocolParser;
use crate::services::csv_recorder::CsvRecorder;
use crate::services::motor_control::{FrameSender, MotorControlService};
use crate::services::realtime_stream::RealtimeStreamService;
use crate::services::session::{MotorSessionService, SessionEvent};
use crate::services::statistics::StatisticsService;
use crate::transport::ble_adapter::BleTransportAdapter;
use crate::ui::main_window::MainWindow;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
pub struct AppController {
pub session: MotorSessionService,
pub stats: StatisticsService,
pub csv: CsvRecorder,
pub current_side: SideKey,
scan_rx: Option<mpsc::Receiver<DeviceCandidate>>,
scan_error_rx: Option<mpsc::Receiver<String>>,
scan_cancel: Option<CancellationToken>,
scan_items: Vec<DeviceCandidate>,
ble_parser: MotorProtocolParser,
usb_parser: MotorProtocolParser,
latest_error: Option<String>,
}
impl AppController {
pub fn new() -> Self {
Self {
session: MotorSessionService::new(),
stats: StatisticsService::new(50),
csv: CsvRecorder::new(),
current_side: SideKey::Right,
scan_rx: None,
scan_error_rx: None,
scan_cancel: None,
scan_items: Vec::new(),
ble_parser: MotorProtocolParser::new(),
usb_parser: MotorProtocolParser::new(),
latest_error: None,
}
}
pub fn build_view_state(&self) -> UiSessionViewState {
UiSessionViewState {
connection_mode: self.session.mode(),
session_state: self.session.state(),
target_label: self.session.target().to_string(),
side_key: Some(self.current_side),
heartbeat_ok: self.session.state() != SessionState::Error,
stream_active: matches!(
self.session.state(),
SessionState::Streaming | SessionState::Recording
),
recording_active: self.csv.recording(),
latest_error: self.latest_error.clone(),
metrics: self.stats.live_metrics(),
}
}
pub fn on_mode_selected(&mut self, mode: ConnectionMode) {
self.stop_scan_task();
self.session.set_mode(mode);
}
pub fn on_side_selected(&mut self, side: SideKey) {
if !matches!(
self.session.state(),
SessionState::Idle | SessionState::Scanning
) {
return;
}
self.current_side = side;
}
pub fn on_scan(&mut self, window: &mut MainWindow) {
self.stop_scan_task();
let (tx, rx) = mpsc::channel(128);
let (error_tx, error_rx) = mpsc::channel(8);
let cancel = CancellationToken::new();
let task_cancel = cancel.clone();
self.scan_rx = Some(rx);
self.scan_error_rx = Some(error_rx);
self.scan_cancel = Some(cancel);
self.scan_items.clear();
self.session.set_scanning();
std::thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build();
if let Ok(rt) = runtime {
if let Err(err) = rt.block_on(BleTransportAdapter::start_scan(tx, task_cancel)) {
let _ = error_tx.blocking_send(format!("BLE scan failed: {err}"));
}
} else {
let _ = error_tx.blocking_send("BLE scan runtime init failed".into());
}
});
window.event_log_panel.append_event("BLE scan started");
}
pub fn on_stop_scan(&mut self, window: &mut MainWindow) {
self.stop_scan_task();
self.session.set_idle_if_scanning();
window.event_log_panel.append_event("BLE scan stopped");
}
pub fn on_connect_ble(&mut self, device_ref: &str, window: &mut MainWindow) {
if device_ref.trim().is_empty() {
self.show_error(window, "请选择 BLE 设备");
return;
}
self.stats.reset();
self.stop_scan_task();
self.session.connect_ble(device_ref, self.current_side);
}
pub fn on_connect_usb(&mut self, port: &str, window: &mut MainWindow) {
if port.trim().is_empty() {
self.show_error(window, "请选择或输入串口");
return;
}
self.stats.reset();
self.session.connect_usb(port, self.current_side);
}
pub fn on_refresh_ports(&mut self, window: &mut MainWindow) {
match serialport::available_ports() {
Ok(ports) => {
let names = ports.into_iter().map(|p| p.port_name).collect::<Vec<_>>();
window.set_ports(names);
window
.event_log_panel
.append_event("Serial ports refreshed");
}
Err(err) => self.show_error(window, &format!("刷新串口失败: {err}")),
}
}
pub fn on_disconnect(&mut self) {
self.stop_scan_task();
self.stop_recording_due_to_session_end(None);
self.session.disconnect();
}
pub fn on_set_params(&mut self, params: serde_json::Value, window: &mut MainWindow) {
let result = {
let mut sender: FrameSender<'_> =
Box::new(|frame| self.session.send_frame(frame).is_ok());
MotorControlService::set_force_mode(&self.current_side, &params, &mut sender)
};
if result.success {
window.event_log_panel.append_event("TX set params command");
} else {
self.show_error(
window,
&result.error.unwrap_or_else(|| "set params failed".into()),
);
}
}
pub fn on_query_current(&mut self, window: &mut MainWindow) {
let result = {
let mut sender: FrameSender<'_> =
Box::new(|frame| self.session.send_frame(frame).is_ok());
MotorControlService::query_current_params(&self.current_side, &mut sender)
};
if result.success {
window
.event_log_panel
.append_event("TX query current params command");
} else {
self.show_error(
window,
&result
.error
.unwrap_or_else(|| "query current failed".into()),
);
}
}
pub fn on_query_mode(&mut self, mode: u8, window: &mut MainWindow) {
let result = {
let mut sender: FrameSender<'_> =
Box::new(|frame| self.session.send_frame(frame).is_ok());
MotorControlService::query_mode_params(&self.current_side, mode, &mut sender)
};
if result.success {
window
.event_log_panel
.append_event(&format!("TX query mode params command: mode={mode}"));
} else {
self.show_error(
window,
&result.error.unwrap_or_else(|| "query mode failed".into()),
);
}
}
pub fn on_toggle_debug(&mut self, enabled: bool, window: &mut MainWindow) {
let result = {
let mut sender: FrameSender<'_> =
Box::new(|frame| self.session.send_frame(frame).is_ok());
if enabled {
MotorControlService::enable_realtime_debug(&self.current_side, &mut sender)
} else {
MotorControlService::disable_realtime_debug(&self.current_side, &mut sender)
}
};
if result.success {
self.session.set_streaming(enabled);
window.event_log_panel.append_event(if enabled {
"TX debug command: enable"
} else {
"TX debug command: disable"
});
} else {
self.show_error(
window,
&result.error.unwrap_or_else(|| {
if enabled {
"enable realtime debug failed".into()
} else {
"disable realtime debug failed".into()
}
}),
);
}
}
pub fn on_start_recording(&mut self, path: &str, window: &mut MainWindow) {
match self.csv.start(std::path::Path::new(path)) {
Ok(()) => {
self.session.set_recording(true);
window
.event_log_panel
.append_event(&format!("Recording started: {path}"));
}
Err(e) => self.show_error(window, &e),
}
}
pub fn on_stop_recording(&mut self, window: &mut MainWindow) {
self.csv.stop();
self.session.set_recording(false);
window.event_log_panel.append_event("Recording stopped");
}
pub fn tick_receive(&mut self, main_window: &mut MainWindow) {
for event in self.session.drain_events() {
match event {
SessionEvent::Ready(msg) => {
self.session.set_connected();
main_window
.event_log_panel
.append_event(&format!("Connected: {msg}"));
}
SessionEvent::Error(error) => {
self.stop_recording_due_to_session_end(Some(main_window));
self.session.set_error();
self.show_error(main_window, &error);
}
}
}
if let Some(rx) = self.scan_rx.as_mut() {
let mut changed = false;
while let Ok(candidate) = rx.try_recv() {
if let Some(existing) = self
.scan_items
.iter_mut()
.find(|item| item.device_ref == candidate.device_ref)
{
*existing = candidate;
} else {
self.scan_items.push(candidate);
}
changed = true;
}
if changed {
main_window.set_scan_items(self.scan_items.clone());
}
}
let mut scan_errors = Vec::new();
if let Some(rx) = self.scan_error_rx.as_mut() {
while let Ok(error) = rx.try_recv() {
scan_errors.push(error);
}
}
for error in scan_errors {
self.show_error(main_window, &error);
}
for frame in self.session.drain_ble_frames() {
let sample = RealtimeStreamService::feed_raw_frame_with_parser(
&mut self.ble_parser,
&frame,
self.session.mode(),
self.session.target(),
self.session.side(),
None,
);
self.consume_sample(sample, main_window);
main_window
.event_log_panel
.append_event(&format!("BLE RX {} bytes", frame.len()));
}
for result in self.session.drain_usb_frames() {
if let Some(error) = result.error {
main_window
.event_log_panel
.append_event(&format!("USB RX error: {error}"));
continue;
}
if let Some(meta) = result.meta {
main_window.event_log_panel.append_event(&format!(
"USB RX meta dev=0x{:02X} len={} seq={} dt={}ms",
meta.dev_id, meta.payload_len, meta.seq_id, meta.dt_ms
));
}
let sample = RealtimeStreamService::feed_raw_frame_with_parser(
&mut self.usb_parser,
&result.payload,
self.session.mode(),
self.session.target(),
self.session.side(),
result.meta,
);
self.consume_sample(sample, main_window);
main_window
.event_log_panel
.append_event(&format!("USB RX {} bytes", result.payload.len()));
}
}
fn consume_sample(
&mut self,
sample: crate::models::MotorDebugSample,
main_window: &mut MainWindow,
) {
self.stats.consume(&sample);
if let Err(err) = self.csv.append(&sample) {
self.show_error(main_window, &format!("CSV write failed: {err}"));
self.csv.stop();
}
main_window.realtime_panel.push_sample(sample);
}
fn show_error(&mut self, window: &mut MainWindow, msg: &str) {
self.latest_error = Some(msg.to_string());
window.show_error(msg);
}
fn stop_scan_task(&mut self) {
if let Some(cancel) = self.scan_cancel.take() {
cancel.cancel();
}
self.scan_rx = None;
self.scan_error_rx = None;
}
fn stop_recording_due_to_session_end(&mut self, window: Option<&mut MainWindow>) {
if self.csv.recording() {
self.csv.stop();
self.session.set_recording(false);
if let Some(window) = window {
window
.event_log_panel
.append_event("Recording stopped: session ended");
}
}
}
}
impl Drop for AppController {
fn drop(&mut self) {
self.stop_scan_task();
self.stop_recording_due_to_session_end(None);
self.session.disconnect();
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::models::{ConnectionMode, UsbUplinkMeta};
use crate::protocol::motor_codec;
#[test]
fn simulated_ble_sample_updates_stats_and_chart() {
let mut controller = AppController::new();
let mut window = MainWindow::default();
let frame = debug_frame(SideKey::Right, "1.0,2.0,3.0,12.0,45.0,0,0,0,10.0");
let sample = RealtimeStreamService::feed_raw_frame_with_parser(
&mut controller.ble_parser,
&frame,
ConnectionMode::Ble,
"ble-test",
SideKey::Right,
None,
);
controller.consume_sample(sample, &mut window);
assert_eq!(controller.stats.live_metrics().sample_count, 1);
assert_eq!(window.realtime_panel.sample_count(), 1);
}
#[test]
fn simulated_usb_sample_preserves_link_metadata() {
let mut controller = AppController::new();
let mut window = MainWindow::default();
let frame = debug_frame(SideKey::Right, "1.0,2.0,3.0,12.0,45.0,0,0,0,10.0");
let sample = RealtimeStreamService::feed_raw_frame_with_parser(
&mut controller.usb_parser,
&frame,
ConnectionMode::UsbDongle,
"COM9",
SideKey::Right,
Some(UsbUplinkMeta {
seq_id: 42,
dt_ms: 20,
dev_id: 0x31,
payload_len: frame.len() + 4,
}),
);
controller.consume_sample(sample, &mut window);
let metrics = controller.stats.live_metrics();
assert_eq!(metrics.sample_count, 1);
assert_eq!(window.realtime_panel.sample_count(), 1);
}
#[test]
fn simulated_bad_debug_sample_surfaces_error_text() {
let mut controller = AppController::new();
let frame = debug_frame(SideKey::Right, "1.0,bad,3.0,12.0,45.0,0,0,0,10.0");
let sample = RealtimeStreamService::feed_raw_frame_with_parser(
&mut controller.ble_parser,
&frame,
ConnectionMode::Ble,
"ble-test",
SideKey::Right,
None,
);
assert_eq!(sample.current, Some(1.0));
assert_eq!(sample.velocity, None);
assert!(sample
.error_text
.as_deref()
.unwrap_or("")
.contains("velocity invalid"));
}
#[test]
fn disconnect_stops_active_recording() {
let path = std::env::temp_dir().join(format!(
"yrobot-panel-disconnect-recording-{}.csv",
std::process::id()
));
let mut controller = AppController::new();
let mut window = MainWindow::default();
controller.on_start_recording(path.to_str().unwrap(), &mut window);
assert!(controller.csv.recording());
controller.on_disconnect();
assert!(!controller.csv.recording());
let _ = std::fs::remove_file(path);
}
#[test]
fn drop_closes_active_recording() {
let path = std::env::temp_dir().join(format!(
"yrobot-panel-drop-recording-{}.csv",
std::process::id()
));
{
let mut controller = AppController::new();
let mut window = MainWindow::default();
controller.on_start_recording(path.to_str().unwrap(), &mut window);
assert!(controller.csv.recording());
}
let content = std::fs::read_to_string(&path).unwrap();
assert!(content.contains("host_ts_ms,connection_mode,target_id"));
let _ = std::fs::remove_file(path);
}
#[test]
fn debug_command_failure_is_shown_to_user() {
let mut controller = AppController::new();
let mut window = MainWindow::default();
controller.on_toggle_debug(true, &mut window);
assert!(window.error_msg.contains("开启实时调试失败"));
assert!(controller
.build_view_state()
.latest_error
.as_deref()
.unwrap_or("")
.contains("开启实时调试失败"));
assert!(window
.event_log_panel
.events()
.iter()
.any(|event| event.contains("ERROR")));
}
fn debug_frame(side: SideKey, text: &str) -> Vec<u8> {
let mut frame = motor_codec::build_debug_toggle(&side, true);
frame[3] = 0x03;
frame.truncate(4);
frame.extend_from_slice(text.as_bytes());
frame[1] = frame.len() as u8;
frame
}
}
+84
View File
@@ -0,0 +1,84 @@
#![cfg_attr(not(debug_assertions), windows_subsystem = "windows")]
mod app;
mod models;
mod protocol;
mod services;
mod transport;
mod ui;
use app::AppController;
use eframe::egui;
use ui::actions::UiAction;
use ui::main_window::MainWindow;
fn main() {
env_logger::init();
let options = eframe::NativeOptions {
viewport: egui::ViewportBuilder::default().with_inner_size([1280.0, 800.0]),
..Default::default()
};
let _ = eframe::run_native(
"YRobot Motor Control Panel",
options,
Box::new(|_cc| Ok(Box::new(YRobotApp::default()))),
);
}
struct YRobotApp {
main_window: MainWindow,
controller: AppController,
}
impl Default for YRobotApp {
fn default() -> Self {
Self {
main_window: MainWindow::default(),
controller: AppController::new(),
}
}
}
impl eframe::App for YRobotApp {
fn update(&mut self, ctx: &egui::Context, _frame: &mut eframe::Frame) {
self.controller.tick_receive(&mut self.main_window);
let vs = self.controller.build_view_state();
self.main_window.update_view_state(vs);
let actions = self.main_window.ui(ctx);
for action in actions {
match action {
UiAction::SelectMode(mode) => self.controller.on_mode_selected(mode),
UiAction::SelectSide(side) => self.controller.on_side_selected(side),
UiAction::Scan => self.controller.on_scan(&mut self.main_window),
UiAction::StopScan => self.controller.on_stop_scan(&mut self.main_window),
UiAction::RefreshPorts => self.controller.on_refresh_ports(&mut self.main_window),
UiAction::ConnectBle(device_ref) => self
.controller
.on_connect_ble(&device_ref, &mut self.main_window),
UiAction::ConnectUsb(port) => {
self.controller.on_connect_usb(&port, &mut self.main_window)
}
UiAction::Disconnect => self.controller.on_disconnect(),
UiAction::SetParams(params) => {
self.controller.on_set_params(params, &mut self.main_window)
}
UiAction::QueryCurrent => self.controller.on_query_current(&mut self.main_window),
UiAction::QueryMode(mode) => {
self.controller.on_query_mode(mode, &mut self.main_window)
}
UiAction::EnableDebug => {
self.controller.on_toggle_debug(true, &mut self.main_window)
}
UiAction::DisableDebug => self
.controller
.on_toggle_debug(false, &mut self.main_window),
UiAction::StartRecording(path) => self
.controller
.on_start_recording(&path, &mut self.main_window),
UiAction::StopRecording => self.controller.on_stop_recording(&mut self.main_window),
}
}
}
}
+160
View File
@@ -0,0 +1,160 @@
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
pub enum ConnectionMode {
#[default]
Ble,
UsbDongle,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
pub enum SessionState {
#[default]
Idle,
Scanning,
Connecting,
Connected,
Streaming,
Recording,
Error,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum SideKey {
Left,
Right,
Neutral,
}
impl Default for SideKey {
fn default() -> Self {
Self::Right
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct UiSessionViewState {
pub connection_mode: ConnectionMode,
pub session_state: SessionState,
pub target_label: String,
pub side_key: Option<SideKey>,
pub heartbeat_ok: bool,
pub stream_active: bool,
pub recording_active: bool,
pub latest_error: Option<String>,
pub metrics: SessionMetrics,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DeviceCandidate {
pub device_name: String,
pub device_ref: String,
pub rssi: Option<i32>,
pub discovered_at_ms: u64,
pub is_preferred: bool,
pub side_hint: Option<SideKey>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub struct UsbUplinkMeta {
pub seq_id: u16,
pub dt_ms: u16,
pub dev_id: u8,
pub payload_len: usize,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MotorDebugSample {
pub host_ts_ms: u64,
pub connection_mode: ConnectionMode,
pub target_id: String,
pub side_key: SideKey,
pub raw_debug_string: String,
pub current: Option<f64>,
pub velocity: Option<f64>,
pub position: Option<f64>,
pub voltage: Option<f64>,
pub motor_temperature: Option<f64>,
pub system_state: Option<String>,
pub error_text: Option<String>,
pub usb_meta: Option<UsbUplinkMeta>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct SessionMetrics {
pub sample_count: u64,
pub dropped_frame_rate: Option<f64>,
pub avg_frame_interval_ms: Option<f64>,
pub max_frame_interval_ms: Option<f64>,
pub avg_usb_dt_ms: Option<f64>,
pub max_usb_dt_ms: Option<f64>,
pub protection_state_count: u64,
pub usb_seq_gap_count: u64,
pub usb_out_of_order_count: u64,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn defaults_match_idle_ble_session() {
assert_eq!(ConnectionMode::default(), ConnectionMode::Ble);
assert_eq!(SessionState::default(), SessionState::Idle);
assert_eq!(SideKey::default(), SideKey::Right);
let view = UiSessionViewState::default();
assert_eq!(view.connection_mode, ConnectionMode::Ble);
assert_eq!(view.session_state, SessionState::Idle);
assert_eq!(view.metrics.sample_count, 0);
assert!(view.latest_error.is_none());
}
#[test]
fn motor_debug_sample_serializes_usb_metadata() {
let sample = MotorDebugSample {
host_ts_ms: 123,
connection_mode: ConnectionMode::UsbDongle,
target_id: "COM9".into(),
side_key: SideKey::Left,
raw_debug_string: "1,2,3,4,5,0,0,0,9".into(),
current: Some(1.0),
velocity: Some(2.0),
position: Some(3.0),
voltage: Some(4.0),
motor_temperature: Some(5.0),
system_state: Some("0".into()),
error_text: None,
usb_meta: Some(UsbUplinkMeta {
seq_id: 7,
dt_ms: 20,
dev_id: 0x31,
payload_len: 32,
}),
};
let json = serde_json::to_string(&sample).unwrap();
let decoded: MotorDebugSample = serde_json::from_str(&json).unwrap();
assert_eq!(decoded.connection_mode, ConnectionMode::UsbDongle);
assert_eq!(decoded.side_key, SideKey::Left);
assert_eq!(decoded.usb_meta.unwrap().seq_id, 7);
}
#[test]
fn session_metrics_serializes_new_usb_dt_fields() {
let metrics = SessionMetrics {
sample_count: 2,
avg_usb_dt_ms: Some(20.0),
max_usb_dt_ms: Some(22.0),
..Default::default()
};
let json = serde_json::to_string(&metrics).unwrap();
let decoded: SessionMetrics = serde_json::from_str(&json).unwrap();
assert_eq!(decoded.sample_count, 2);
assert_eq!(decoded.avg_usb_dt_ms, Some(20.0));
assert_eq!(decoded.max_usb_dt_ms, Some(22.0));
}
}
+3
View File
@@ -0,0 +1,3 @@
pub mod motor_codec;
pub mod motor_parser;
pub mod usb_codec;
+81
View File
@@ -0,0 +1,81 @@
use crate::models::SideKey;
use serde_json::json;
pub fn side_key_for_byte(b: u8) -> Option<SideKey> {
match b {
0x6F | 0x6E => Some(SideKey::Left),
0xAF | 0xAE => Some(SideKey::Right),
0x2F | 0x2E => Some(SideKey::Neutral),
_ => None,
}
}
const CMD_HEARTBEAT: u8 = 0x00;
const CMD_SET_PARAMS: u8 = 0x01;
const CMD_QUERY_PARAMS: u8 = 0x02;
const CMD_DEBUG_TOGGLE: u8 = 0x03;
fn side_byte(side: &SideKey) -> u8 {
match side {
SideKey::Left => 0x6F,
SideKey::Right => 0xAF,
SideKey::Neutral => 0x2F,
}
}
fn build_frame(side: &SideKey, command: u8, payload: &[u8]) -> Vec<u8> {
let inner: Vec<u8> = std::iter::once(command)
.chain(payload.iter().copied())
.collect();
let length: u8 = (inner.len() + 3)
.try_into()
.expect("motor frame length is validated by callers");
let mut frame = vec![side_byte(side), length, 0x00];
frame.extend(inner);
frame
}
pub fn build_heartbeat(side: &SideKey) -> Vec<u8> {
build_frame(side, CMD_HEARTBEAT, &[])
}
pub fn build_set_params(side: &SideKey, params: &serde_json::Value) -> Result<Vec<u8>, String> {
let payload = serde_json::to_vec(params).map_err(|err| err.to_string())?;
if payload.len() > 251 {
return Err("参数 JSON 过长,无法编码到单个电机帧".into());
}
Ok(build_frame(side, CMD_SET_PARAMS, &payload))
}
pub fn build_query_current(side: &SideKey) -> Vec<u8> {
build_frame(side, CMD_QUERY_PARAMS, &[])
}
pub fn build_query_mode(side: &SideKey, mode: u8) -> Vec<u8> {
build_frame(
side,
CMD_QUERY_PARAMS,
&serde_json::to_vec(&json!({"fm_mode": mode})).expect("query mode JSON is built from a u8"),
)
}
pub fn build_debug_toggle(side: &SideKey, enabled: bool) -> Vec<u8> {
build_frame(
side,
CMD_DEBUG_TOGGLE,
if enabled { b"enable" } else { b"disable" },
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn heartbeat_frame_has_expected_header() {
let frame = build_heartbeat(&SideKey::Right);
assert_eq!(frame[0], 0xAF);
assert_eq!(frame[1], 4);
assert_eq!(frame[3], CMD_HEARTBEAT);
}
}
+357
View File
@@ -0,0 +1,357 @@
use crate::models::SideKey;
use std::collections::HashMap;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MotorEventKind {
TuningParams,
DebugSample,
Unknown,
}
#[derive(Debug, Clone)]
pub struct MotorProtocolEvent {
pub kind: MotorEventKind,
pub side_key: Option<SideKey>,
pub params_json: Option<String>,
pub params_dict: Option<serde_json::Value>,
pub params_error: Option<String>,
pub raw_debug_string: Option<String>,
pub debug_current: Option<f64>,
pub debug_velocity: Option<f64>,
pub debug_position: Option<f64>,
pub debug_voltage: Option<f64>,
pub debug_temperature: Option<f64>,
pub debug_sys: Option<i32>,
pub debug_merr: Option<i32>,
pub debug_obd: Option<i32>,
pub debug_device_time: Option<f64>,
pub debug_parse_error: Option<String>,
pub raw_bytes: Vec<u8>,
}
impl MotorProtocolEvent {
fn new(kind: MotorEventKind) -> Self {
Self {
kind,
side_key: None,
params_json: None,
params_dict: None,
params_error: None,
raw_debug_string: None,
debug_current: None,
debug_velocity: None,
debug_position: None,
debug_voltage: None,
debug_temperature: None,
debug_sys: None,
debug_merr: None,
debug_obd: None,
debug_device_time: None,
debug_parse_error: None,
raw_bytes: vec![],
}
}
}
pub struct MotorProtocolParser {
buf: Vec<u8>,
pending_params: HashMap<u8, PendingParamPackets>,
}
impl MotorProtocolParser {
pub fn new() -> Self {
Self {
buf: vec![],
pending_params: HashMap::new(),
}
}
pub fn feed(&mut self, frame: &[u8]) -> Option<MotorProtocolEvent> {
self.buf.extend_from_slice(frame);
if self.buf.len() < 4 {
return None;
}
let frame_len = self.buf[1] as usize;
if frame_len < 4 {
let raw = std::mem::take(&mut self.buf);
let mut ev = MotorProtocolEvent::new(MotorEventKind::Unknown);
ev.raw_bytes = raw;
ev.debug_parse_error = Some("invalid motor frame length".into());
return Some(ev);
}
if self.buf.len() < frame_len {
return None;
}
let packet = self.buf.drain(..frame_len).collect::<Vec<_>>();
let kind = match packet[3] {
0x02 => MotorEventKind::TuningParams,
0x03 => MotorEventKind::DebugSample,
_ => MotorEventKind::Unknown,
};
let mut ev = MotorProtocolEvent::new(kind);
ev.raw_bytes = packet.clone();
ev.side_key = super::motor_codec::side_key_for_byte(packet[0]);
match kind {
MotorEventKind::TuningParams => {
let payload = &packet[4..];
let text = String::from_utf8_lossy(payload).to_string();
match self.parse_tuning_params(packet[0], &text) {
ParamParseOutcome::Ready(json_text, value, error) => {
ev.params_json = Some(json_text);
ev.params_dict = value;
ev.params_error = error;
}
ParamParseOutcome::Waiting => return None,
}
}
MotorEventKind::DebugSample => {
let text = String::from_utf8_lossy(&packet[4..]).to_string();
ev.raw_debug_string = Some(text.clone());
let fields = text.split(',').map(str::trim).collect::<Vec<_>>();
let mut errors = Vec::new();
ev.debug_current = parse_f64_field(&fields, 0, "current", &mut errors);
ev.debug_velocity = parse_f64_field(&fields, 1, "velocity", &mut errors);
ev.debug_position = parse_f64_field(&fields, 2, "position", &mut errors);
ev.debug_voltage = parse_f64_field(&fields, 3, "voltage", &mut errors);
ev.debug_temperature = parse_f64_field(&fields, 4, "temperature", &mut errors);
ev.debug_sys = parse_i32_field(&fields, 5, "sys", &mut errors);
ev.debug_merr = parse_i32_field(&fields, 6, "merr", &mut errors);
ev.debug_obd = parse_i32_field(&fields, 7, "obd", &mut errors);
ev.debug_device_time = parse_f64_field(&fields, 8, "device_time", &mut errors);
if !errors.is_empty() {
ev.debug_parse_error =
Some(format!("debug field parse error: {}", errors.join(", ")));
}
}
MotorEventKind::Unknown => {}
}
Some(ev)
}
fn parse_tuning_params(&mut self, side_byte: u8, text: &str) -> ParamParseOutcome {
let json = match serde_json::from_str::<serde_json::Value>(text) {
Ok(json) => json,
Err(err) => {
return ParamParseOutcome::Ready(text.to_string(), None, Some(err.to_string()));
}
};
let Some(count) = json.get("packet_count").and_then(|value| value.as_u64()) else {
return ParamParseOutcome::Ready(text.to_string(), Some(json), None);
};
let Some(index) = json.get("packet_index").and_then(|value| value.as_u64()) else {
return ParamParseOutcome::Ready(
text.to_string(),
None,
Some("missing packet_index".into()),
);
};
if count == 0 || count > 64 || index >= count {
return ParamParseOutcome::Ready(
text.to_string(),
None,
Some("invalid packet_index/packet_count".into()),
);
}
let Some(chunk) = json
.get("data")
.or_else(|| json.get("payload"))
.or_else(|| json.get("json"))
.or_else(|| json.get("chunk"))
.and_then(|value| value.as_str())
else {
return ParamParseOutcome::Ready(
text.to_string(),
None,
Some("missing packet data".into()),
);
};
let pending = self
.pending_params
.entry(side_byte)
.or_insert_with(|| PendingParamPackets::new(count as usize));
if pending.packet_count != count as usize {
// 分包总数变化意味着上一轮响应已经失配,重建缓存避免混包。
*pending = PendingParamPackets::new(count as usize);
}
pending.parts[index as usize] = Some(chunk.to_string());
if pending.parts.iter().any(Option::is_none) {
return ParamParseOutcome::Waiting;
}
let Some(pending) = self.pending_params.remove(&side_byte) else {
return ParamParseOutcome::Ready(
text.to_string(),
None,
Some("missing parameter packet cache".into()),
);
};
let full_json = pending
.parts
.into_iter()
.map(|part| part.unwrap_or_default())
.collect::<String>();
match serde_json::from_str::<serde_json::Value>(&full_json) {
Ok(value) => ParamParseOutcome::Ready(full_json, Some(value), None),
Err(err) => ParamParseOutcome::Ready(full_json, None, Some(err.to_string())),
}
}
}
fn parse_f64_field(
fields: &[&str],
index: usize,
name: &str,
errors: &mut Vec<String>,
) -> Option<f64> {
let Some(value) = fields.get(index).copied() else {
errors.push(format!("{name} missing"));
return None;
};
match value.parse::<f64>() {
Ok(parsed) if parsed.is_finite() => Some(parsed),
_ => {
errors.push(format!("{name} invalid"));
None
}
}
}
fn parse_i32_field(
fields: &[&str],
index: usize,
name: &str,
errors: &mut Vec<String>,
) -> Option<i32> {
let Some(value) = fields.get(index).copied() else {
errors.push(format!("{name} missing"));
return None;
};
match value.parse::<i32>() {
Ok(parsed) => Some(parsed),
Err(_) => {
errors.push(format!("{name} invalid"));
None
}
}
}
struct PendingParamPackets {
packet_count: usize,
parts: Vec<Option<String>>,
}
impl PendingParamPackets {
fn new(packet_count: usize) -> Self {
Self {
packet_count,
parts: vec![None; packet_count],
}
}
}
enum ParamParseOutcome {
Ready(String, Option<serde_json::Value>, Option<String>),
Waiting,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_debug_sample_fields() {
let mut parser = MotorProtocolParser::new();
let mut frame = b"\xAF\x00\x00\x031.5,2.5,3.5,12.0,45.0,0,1,2,99.0".to_vec();
frame[1] = frame.len() as u8;
let event = parser.feed(&frame).unwrap();
assert_eq!(event.kind, MotorEventKind::DebugSample);
assert_eq!(event.side_key, Some(SideKey::Right));
assert_eq!(event.debug_current, Some(1.5));
assert_eq!(event.debug_velocity, Some(2.5));
assert_eq!(event.debug_position, Some(3.5));
assert_eq!(event.debug_voltage, Some(12.0));
assert_eq!(event.debug_temperature, Some(45.0));
assert_eq!(event.debug_sys, Some(0));
assert_eq!(event.debug_merr, Some(1));
assert_eq!(event.debug_obd, Some(2));
assert_eq!(event.debug_device_time, Some(99.0));
}
#[test]
fn marks_bad_debug_fields_without_blocking_valid_fields() {
let mut parser = MotorProtocolParser::new();
let mut frame = b"\xAF\x00\x00\x031.5,bad,3.5,12.0,45.0,0,1".to_vec();
frame[1] = frame.len() as u8;
let event = parser.feed(&frame).unwrap();
assert_eq!(event.kind, MotorEventKind::DebugSample);
assert_eq!(event.debug_current, Some(1.5));
assert_eq!(event.debug_velocity, None);
assert_eq!(event.debug_position, Some(3.5));
let error = event.debug_parse_error.unwrap();
assert!(error.contains("velocity invalid"));
assert!(error.contains("obd missing"));
assert!(error.contains("device_time missing"));
}
#[test]
fn parses_tuning_params_json() {
let mut parser = MotorProtocolParser::new();
let mut frame = b"\x6F\x00\x00\x02{\"fm_mode\":2}".to_vec();
frame[1] = frame.len() as u8;
let event = parser.feed(&frame).unwrap();
assert_eq!(event.kind, MotorEventKind::TuningParams);
assert_eq!(event.side_key, Some(SideKey::Left));
assert_eq!(event.params_dict.unwrap()["fm_mode"], 2);
assert!(event.params_error.is_none());
}
#[test]
fn reassembles_tuning_params_packets() {
let mut parser = MotorProtocolParser::new();
let first = tuning_frame(r#"{"packet_index":0,"packet_count":2,"data":"{\"fm_mode\":"}"#);
let second =
tuning_frame(r#"{"packet_index":1,"packet_count":2,"data":"2,\"fm_mset\":3.5}"}"#);
assert!(parser.feed(&first).is_none());
let event = parser.feed(&second).unwrap();
assert_eq!(event.kind, MotorEventKind::TuningParams);
assert_eq!(event.params_dict.unwrap()["fm_mset"], 3.5);
assert!(event.params_error.is_none());
}
#[test]
fn reports_missing_tuning_packet_data() {
let mut parser = MotorProtocolParser::new();
let frame = tuning_frame(r#"{"packet_index":0,"packet_count":2}"#);
let event = parser.feed(&frame).unwrap();
assert_eq!(event.kind, MotorEventKind::TuningParams);
assert!(event.params_error.unwrap().contains("missing packet data"));
}
#[test]
fn waits_for_complete_frame() {
let mut parser = MotorProtocolParser::new();
assert!(parser.feed(b"\xAF\x08\x00\x031.").is_none());
let event = parser.feed(b"5,2").unwrap();
assert_eq!(event.kind, MotorEventKind::DebugSample);
assert_eq!(event.debug_current, Some(1.5));
}
fn tuning_frame(payload: &str) -> Vec<u8> {
let mut frame = Vec::from([0x6F, 0x00, 0x00, 0x02]);
frame.extend_from_slice(payload.as_bytes());
frame[1] = frame.len() as u8;
frame
}
}
+170
View File
@@ -0,0 +1,170 @@
use crate::models::UsbUplinkMeta;
use crc::Crc;
use std::collections::VecDeque;
const CRC_8_MAXIM: crc::Algorithm<u8> = crc::Algorithm {
width: 8,
poly: 0x31,
init: 0x00,
refin: true,
refout: true,
xorout: 0x00,
check: 0xA1,
residue: 0x00,
};
#[derive(Debug, Clone, Default)]
pub struct UsbUplinkResult {
pub meta: Option<UsbUplinkMeta>,
pub payload: Vec<u8>,
pub error: Option<String>,
}
#[derive(Debug, Default)]
pub struct UsbStreamAssembler {
buf: VecDeque<u8>,
}
pub fn encode_downlink(motor_payload: &[u8]) -> Vec<u8> {
let mut out = vec![0x5A, 0xA5, 0x31, motor_payload.len() as u8];
out.extend_from_slice(motor_payload);
let crc = Crc::<u8>::new(&CRC_8_MAXIM).checksum(&out);
out.push(crc);
out
}
pub fn decode_uplink(frame: &[u8]) -> Result<UsbUplinkResult, String> {
if frame.len() < 7 {
return Err("USB frame too short".into());
}
if frame[0] != 0x5A || frame[1] != 0xA5 {
return Err("USB sync mismatch".into());
}
let crc = Crc::<u8>::new(&CRC_8_MAXIM).checksum(&frame[..frame.len() - 1]);
if crc != frame[frame.len() - 1] {
return Err("USB CRC mismatch".into());
}
let dev = frame[2];
let len = frame[3] as usize;
if frame.len() != len + 5 {
return Err("USB length mismatch".into());
}
let payload = frame[4..4 + len].to_vec();
if payload.len() < 4 {
return Err("USB payload too short".into());
}
let seq_id = u16::from_le_bytes([payload[0], payload[1]]);
let dt_ms = u16::from_le_bytes([payload[2], payload[3]]);
Ok(UsbUplinkResult {
meta: Some(UsbUplinkMeta {
seq_id,
dt_ms,
dev_id: dev,
payload_len: len,
}),
payload: payload[4..].to_vec(),
error: None,
})
}
impl UsbStreamAssembler {
pub fn push(&mut self, bytes: &[u8]) -> Vec<UsbUplinkResult> {
self.buf.extend(bytes.iter().copied());
let mut out = Vec::new();
loop {
while self.buf.len() >= 2 {
if self.buf[0] == 0x5A && self.buf[1] == 0xA5 {
break;
}
self.buf.pop_front();
}
if self.buf.len() < 5 {
break;
}
let len = self.buf[3] as usize;
let frame_len = len + 5;
if self.buf.len() < frame_len {
break;
}
let frame: Vec<u8> = self.buf.drain(..frame_len).collect();
match decode_uplink(&frame) {
Ok(decoded) => out.push(decoded),
Err(err) => out.push(UsbUplinkResult {
meta: None,
payload: vec![],
error: Some(err),
}),
}
}
out
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn crc8_maxim_uses_standard_check_vector() {
let crc = Crc::<u8>::new(&CRC_8_MAXIM).checksum(b"123456789");
assert_eq!(crc, 0xA1);
}
#[test]
fn encode_then_decode_round_trip() {
let payload = vec![1, 2, 3, 4, 5, 6];
let frame = encode_downlink(&payload);
let decoded = decode_uplink(&frame).unwrap();
let meta = decoded.meta.unwrap();
assert_eq!(meta.seq_id, 0x0201);
assert_eq!(meta.dt_ms, 0x0403);
assert_eq!(decoded.payload, vec![5, 6]);
}
#[test]
fn stream_assembler_waits_for_complete_frame() {
let frame = encode_downlink(&[1, 0, 20, 0, 0xAF, 4, 0, 0]);
let mut assembler = UsbStreamAssembler::default();
assert!(assembler.push(&frame[..3]).is_empty());
let out = assembler.push(&frame[3..]);
assert_eq!(out.len(), 1);
assert_eq!(out[0].meta.unwrap().seq_id, 1);
assert_eq!(out[0].payload, vec![0xAF, 4, 0, 0]);
}
#[test]
fn stream_assembler_skips_noise_before_sync() {
let frame = encode_downlink(&[2, 0, 20, 0, 0xAF, 4, 0, 0]);
let mut bytes = vec![0x00, 0x11, 0x5A, 0x00];
bytes.extend(frame);
let mut assembler = UsbStreamAssembler::default();
let out = assembler.push(&bytes);
assert_eq!(out.len(), 1);
assert_eq!(out[0].meta.unwrap().seq_id, 2);
}
#[test]
fn stream_assembler_reports_crc_error_and_recovers() {
let mut bad = encode_downlink(&[3, 0, 20, 0, 0xAF, 4, 0, 0]);
let good = encode_downlink(&[4, 0, 20, 0, 0xAF, 4, 0, 0]);
let last = bad.len() - 1;
bad[last] ^= 0xFF;
bad.extend(good);
let mut assembler = UsbStreamAssembler::default();
let out = assembler.push(&bad);
assert_eq!(out.len(), 2);
assert!(out[0].error.as_deref().unwrap_or("").contains("CRC"));
assert_eq!(out[1].meta.unwrap().seq_id, 4);
}
}
+144
View File
@@ -0,0 +1,144 @@
use crate::models::MotorDebugSample;
use csv::Writer;
use std::fs::File;
use std::path::Path;
#[derive(Default)]
pub struct CsvRecorder {
recording: bool,
writer: Option<Writer<File>>,
}
impl CsvRecorder {
pub fn new() -> Self {
Self::default()
}
pub fn start(&mut self, _path: &Path) -> Result<(), String> {
let file = File::create(_path).map_err(|e| e.to_string())?;
let mut writer = Writer::from_writer(file);
writer
.write_record([
"host_ts_ms",
"connection_mode",
"target_id",
"side_key",
"raw_debug_string",
"current",
"velocity",
"position",
"voltage",
"motor_temperature",
"system_state",
"error_text",
"usb_seq_id",
"usb_dt_ms",
])
.map_err(|e| e.to_string())?;
self.writer = Some(writer);
self.recording = true;
Ok(())
}
pub fn stop(&mut self) {
if let Some(mut writer) = self.writer.take() {
let _ = writer.flush();
}
self.recording = false;
}
pub fn recording(&self) -> bool {
self.recording
}
pub fn append(&mut self, sample: &MotorDebugSample) -> Result<(), String> {
if let Some(writer) = self.writer.as_mut() {
writer
.serialize((
sample.host_ts_ms,
format!("{:?}", sample.connection_mode),
sample.target_id.clone(),
format!("{:?}", sample.side_key),
sample.raw_debug_string.clone(),
sample.current,
sample.velocity,
sample.position,
sample.voltage,
sample.motor_temperature,
sample.system_state.clone(),
sample.error_text.clone(),
sample.usb_meta.map(|m| m.seq_id),
sample.usb_meta.map(|m| m.dt_ms),
))
.map_err(|e| e.to_string())?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::models::{ConnectionMode, MotorDebugSample, SideKey, UsbUplinkMeta};
#[test]
fn writes_header_and_sample() {
let path =
std::env::temp_dir().join(format!("yrobot-panel-csv-test-{}.csv", std::process::id()));
let mut recorder = CsvRecorder::new();
recorder.start(&path).unwrap();
recorder.append(&sample()).unwrap();
recorder.stop();
let content = std::fs::read_to_string(&path).unwrap();
let _ = std::fs::remove_file(&path);
assert!(content.contains("host_ts_ms,connection_mode,target_id"));
assert!(content.contains("123,UsbDongle,COM3"));
assert!(content.contains(",9,20"));
}
#[test]
fn append_without_recording_is_noop() {
let mut recorder = CsvRecorder::new();
let sample = sample();
assert!(recorder.append(&sample).is_ok());
assert!(!recorder.recording());
}
#[test]
fn start_reports_path_errors() {
let mut recorder = CsvRecorder::new();
let path = std::env::temp_dir()
.join(format!("missing-yrobot-dir-{}", std::process::id()))
.join("out.csv");
let result = recorder.start(&path);
assert!(result.is_err());
assert!(!recorder.recording());
}
fn sample() -> MotorDebugSample {
MotorDebugSample {
host_ts_ms: 123,
connection_mode: ConnectionMode::UsbDongle,
target_id: "COM3".into(),
side_key: SideKey::Right,
raw_debug_string: "raw".into(),
current: Some(1.0),
velocity: Some(2.0),
position: Some(3.0),
voltage: Some(4.0),
motor_temperature: Some(5.0),
system_state: Some("0".into()),
error_text: None,
usb_meta: Some(UsbUplinkMeta {
seq_id: 9,
dt_ms: 20,
dev_id: 0x31,
payload_len: 8,
}),
}
}
}
+142
View File
@@ -0,0 +1,142 @@
use crate::models::SideKey;
use crate::protocol::motor_codec;
use crate::services::session::SessionEvent;
use std::any::Any;
use std::panic::{catch_unwind, AssertUnwindSafe};
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::Duration;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
pub type FrameWriter = dyn FnMut(&[u8]) -> bool + Send;
pub struct HeartbeatService;
impl HeartbeatService {
pub fn start(
side: SideKey,
_mode: crate::models::ConnectionMode,
writer: Arc<Mutex<Box<FrameWriter>>>,
cancel: CancellationToken,
event_tx: Option<mpsc::Sender<SessionEvent>>,
) -> thread::JoinHandle<()> {
thread::spawn(move || {
let result = catch_unwind(AssertUnwindSafe(|| {
while !cancel.is_cancelled() {
let heartbeat = motor_codec::build_heartbeat(&side);
match writer.lock() {
Ok(mut guard) => {
if !guard(&heartbeat) {
send_error(&event_tx, "heartbeat send failed".into());
}
}
Err(_) => send_error(&event_tx, "heartbeat writer lock poisoned".into()),
}
if wait_or_cancel(&cancel, Duration::from_secs(2)) {
break;
}
}
}));
if let Err(payload) = result {
send_error(
&event_tx,
format!("heartbeat task panic: {}", panic_message(payload)),
);
}
})
}
}
fn wait_or_cancel(cancel: &CancellationToken, duration: Duration) -> bool {
let step = Duration::from_millis(50);
let start = std::time::Instant::now();
while start.elapsed() < duration {
if cancel.is_cancelled() {
return true;
}
thread::sleep(step.min(duration.saturating_sub(start.elapsed())));
}
cancel.is_cancelled()
}
fn send_error(event_tx: &Option<mpsc::Sender<SessionEvent>>, error: String) {
if let Some(tx) = event_tx {
let _ = tx.blocking_send(SessionEvent::Error(error));
}
}
fn panic_message(payload: Box<dyn Any + Send>) -> String {
if let Some(message) = payload.downcast_ref::<&str>() {
(*message).to_string()
} else if let Some(message) = payload.downcast_ref::<String>() {
message.clone()
} else {
"unknown panic".into()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::models::{ConnectionMode, SideKey};
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn heartbeat_stops_promptly_after_cancel() {
let count = Arc::new(AtomicUsize::new(0));
let writer_count = count.clone();
let writer: Arc<Mutex<Box<FrameWriter>>> = Arc::new(Mutex::new(Box::new(move |_frame| {
writer_count.fetch_add(1, Ordering::SeqCst);
true
})));
let cancel = CancellationToken::new();
let handle = HeartbeatService::start(
SideKey::Right,
ConnectionMode::Ble,
writer,
cancel.clone(),
None,
);
wait_until(|| count.load(Ordering::SeqCst) >= 1);
cancel.cancel();
let after_cancel = count.load(Ordering::SeqCst);
std::thread::sleep(Duration::from_millis(200));
assert_eq!(count.load(Ordering::SeqCst), after_cancel);
handle.join().unwrap();
}
#[test]
fn heartbeat_reports_writer_failure() {
let writer: Arc<Mutex<Box<FrameWriter>>> =
Arc::new(Mutex::new(Box::new(move |_frame| false)));
let cancel = CancellationToken::new();
let (tx, mut rx) = mpsc::channel(4);
let handle = HeartbeatService::start(
SideKey::Right,
ConnectionMode::Ble,
writer,
cancel.clone(),
Some(tx),
);
let event = rx.blocking_recv().unwrap();
cancel.cancel();
handle.join().unwrap();
match event {
SessionEvent::Error(error) => assert!(error.contains("heartbeat send failed")),
_ => panic!("unexpected heartbeat event"),
}
}
fn wait_until(predicate: impl Fn() -> bool) {
let start = std::time::Instant::now();
while !predicate() {
assert!(start.elapsed() < Duration::from_secs(1));
std::thread::sleep(Duration::from_millis(10));
}
}
}
+6
View File
@@ -0,0 +1,6 @@
pub mod csv_recorder;
pub mod heartbeat;
pub mod motor_control;
pub mod realtime_stream;
pub mod session;
pub mod statistics;
+227
View File
@@ -0,0 +1,227 @@
use crate::models::SideKey;
use crate::protocol::motor_codec;
pub type FrameSender<'a> = Box<dyn FnMut(&[u8]) -> bool + Send + 'a>;
pub struct MotorCommandResult {
pub success: bool,
pub error: Option<String>,
}
pub struct MotorControlService;
impl MotorControlService {
pub fn set_force_mode<'a>(
side: &SideKey,
params: &serde_json::Value,
sender: &mut FrameSender<'a>,
) -> MotorCommandResult {
if let Err(error) = Self::validate_force_params(params) {
return MotorCommandResult {
success: false,
error: Some(error),
};
}
let frame = match motor_codec::build_set_params(side, params) {
Ok(frame) => frame,
Err(error) => {
return MotorCommandResult {
success: false,
error: Some(error),
}
}
};
if sender(&frame) {
MotorCommandResult {
success: true,
error: None,
}
} else {
MotorCommandResult {
success: false,
error: Some("发送参数帧失败".into()),
}
}
}
pub fn query_current_params<'a>(
side: &SideKey,
sender: &mut FrameSender<'a>,
) -> MotorCommandResult {
let frame = motor_codec::build_query_current(side);
if sender(&frame) {
MotorCommandResult {
success: true,
error: None,
}
} else {
MotorCommandResult {
success: false,
error: Some("查询当前参数失败".into()),
}
}
}
pub fn query_mode_params<'a>(
side: &SideKey,
mode: u8,
sender: &mut FrameSender<'a>,
) -> MotorCommandResult {
if mode > 6 {
return MotorCommandResult {
success: false,
error: Some("力控模式必须在 0..6".into()),
};
}
let frame = motor_codec::build_query_mode(side, mode);
if sender(&frame) {
MotorCommandResult {
success: true,
error: None,
}
} else {
MotorCommandResult {
success: false,
error: Some("查询模式参数失败".into()),
}
}
}
pub fn enable_realtime_debug<'a>(
side: &SideKey,
sender: &mut FrameSender<'a>,
) -> MotorCommandResult {
let frame = motor_codec::build_debug_toggle(side, true);
if sender(&frame) {
MotorCommandResult {
success: true,
error: None,
}
} else {
MotorCommandResult {
success: false,
error: Some("开启实时调试失败".into()),
}
}
}
pub fn disable_realtime_debug<'a>(
side: &SideKey,
sender: &mut FrameSender<'a>,
) -> MotorCommandResult {
let frame = motor_codec::build_debug_toggle(side, false);
if sender(&frame) {
MotorCommandResult {
success: true,
error: None,
}
} else {
MotorCommandResult {
success: false,
error: Some("关闭实时调试失败".into()),
}
}
}
fn validate_force_params(params: &serde_json::Value) -> Result<(), String> {
let object = params
.as_object()
.ok_or_else(|| "参数必须是 JSON 对象".to_string())?;
if object.is_empty() {
return Err("参数不能为空".into());
}
let mode = object
.get("fm_mode")
.and_then(|value| value.as_u64())
.ok_or_else(|| "缺少 fm_mode 或 fm_mode 不是整数".to_string())?;
if mode > 6 {
return Err("fm_mode 必须在 0..6".into());
}
for (key, value) in object {
if key == "fm_mode" {
continue;
}
let number = value.as_f64().ok_or_else(|| format!("{key} 必须是数字"))?;
if !number.is_finite() || number < 0.0 {
return Err(format!("{key} 必须是非负有限数"));
}
}
let payload = serde_json::to_vec(params).map_err(|err| err.to_string())?;
if payload.len() > 251 {
// 电机帧 length 是 u8,提前拒绝可以避免编码层 unwrap panic。
return Err("参数 JSON 过长,无法编码到单个电机帧".into());
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
fn capture_sender<'a>(sent: &'a mut Vec<Vec<u8>>) -> FrameSender<'a> {
Box::new(move |frame| {
sent.push(frame.to_vec());
true
})
}
#[test]
fn sends_query_current_frame() {
let mut sent = Vec::new();
let mut sender = capture_sender(&mut sent);
let result = MotorControlService::query_current_params(&SideKey::Neutral, &mut sender);
assert!(result.success);
drop(sender);
assert_eq!(sent[0][0], 0x2F);
assert_eq!(sent[0][3], 0x02);
}
#[test]
fn sends_force_params_from_payload() {
let mut sent = Vec::new();
let mut sender = capture_sender(&mut sent);
let params = serde_json::json!({"fm_mode": 4, "fm_mset": 1.2, "fm_vset": 3.4});
let result = MotorControlService::set_force_mode(&SideKey::Right, &params, &mut sender);
assert!(result.success);
drop(sender);
assert_eq!(sent[0][0], 0xAF);
assert_eq!(sent[0][3], 0x01);
let payload = std::str::from_utf8(&sent[0][4..]).unwrap();
assert!(payload.contains("\"fm_mode\":4"));
assert!(payload.contains("\"fm_mset\":1.2"));
}
#[test]
fn rejects_invalid_force_params() {
let mut sent = Vec::new();
let mut sender = capture_sender(&mut sent);
let params = serde_json::json!({"fm_mode": 7, "fm_mset": -1.0});
let result = MotorControlService::set_force_mode(&SideKey::Right, &params, &mut sender);
assert!(!result.success);
drop(sender);
assert!(sent.is_empty());
}
#[test]
fn rejects_invalid_query_mode() {
let mut sent = Vec::new();
let mut sender = capture_sender(&mut sent);
let result = MotorControlService::query_mode_params(&SideKey::Right, 7, &mut sender);
assert!(!result.success);
drop(sender);
assert!(sent.is_empty());
}
#[test]
fn sends_debug_toggle_frame() {
let mut sent = Vec::new();
let mut sender = capture_sender(&mut sent);
let result = MotorControlService::enable_realtime_debug(&SideKey::Left, &mut sender);
assert!(result.success);
drop(sender);
assert_eq!(sent[0][0], 0x6F);
assert_eq!(sent[0][3], 0x03);
assert!(sent[0].ends_with(b"enable"));
}
}
+55
View File
@@ -0,0 +1,55 @@
use crate::models::{ConnectionMode, MotorDebugSample, SideKey, UsbUplinkMeta};
use crate::protocol::motor_parser::{MotorEventKind, MotorProtocolParser};
pub struct RealtimeStreamService;
impl RealtimeStreamService {
pub fn feed_raw_frame_with_parser(
parser: &mut MotorProtocolParser,
frame: &[u8],
mode: ConnectionMode,
target: &str,
side: SideKey,
usb_meta: Option<UsbUplinkMeta>,
) -> MotorDebugSample {
let event = parser.feed(frame);
let event_kind = event.as_ref().map(|ev| ev.kind);
let raw_debug_string = event
.as_ref()
.and_then(|ev| ev.raw_debug_string.clone())
.unwrap_or_else(|| String::from_utf8_lossy(frame).to_string());
MotorDebugSample {
host_ts_ms: current_time_ms(),
connection_mode: mode,
target_id: target.to_string(),
side_key: side,
raw_debug_string,
current: event.as_ref().and_then(|ev| ev.debug_current),
velocity: event.as_ref().and_then(|ev| ev.debug_velocity),
position: event.as_ref().and_then(|ev| ev.debug_position),
voltage: event.as_ref().and_then(|ev| ev.debug_voltage),
motor_temperature: event.as_ref().and_then(|ev| ev.debug_temperature),
system_state: event
.as_ref()
.and_then(|ev| ev.debug_sys.map(|v| v.to_string())),
error_text: event.as_ref().and_then(|ev| {
ev.debug_parse_error
.clone()
.or(ev.params_error.clone())
.or_else(|| match event_kind {
Some(MotorEventKind::Unknown) => Some("unknown motor event".into()),
_ => None,
})
}),
usb_meta,
}
}
}
fn current_time_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
+291
View File
@@ -0,0 +1,291 @@
use crate::models::{ConnectionMode, SessionState, SideKey};
use crate::protocol::usb_codec;
use crate::services::heartbeat::{FrameWriter, HeartbeatService};
use crate::transport::ble_adapter::BleTransportAdapter;
use crate::transport::serial_adapter::SerialTransportAdapter;
use std::any::Any;
use std::panic::{catch_unwind, AssertUnwindSafe};
use std::sync::{Arc, Mutex};
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
pub struct MotorSessionService {
mode: ConnectionMode,
state: SessionState,
target: String,
side: SideKey,
cancel: Option<CancellationToken>,
writer: Option<Arc<Mutex<Box<FrameWriter>>>>,
ble_rx: Option<mpsc::Receiver<Vec<u8>>>,
usb_rx: Option<mpsc::Receiver<usb_codec::UsbUplinkResult>>,
event_rx: Option<mpsc::Receiver<SessionEvent>>,
}
#[derive(Debug, Clone)]
pub enum SessionEvent {
Ready(String),
Error(String),
}
impl MotorSessionService {
pub fn new() -> Self {
Self {
mode: ConnectionMode::Ble,
state: SessionState::Idle,
target: "none".into(),
side: SideKey::Right,
cancel: None,
writer: None,
ble_rx: None,
usb_rx: None,
event_rx: None,
}
}
pub fn mode(&self) -> ConnectionMode {
self.mode
}
pub fn state(&self) -> SessionState {
self.state
}
pub fn target(&self) -> &str {
&self.target
}
pub fn side(&self) -> SideKey {
self.side
}
pub fn set_mode(&mut self, mode: ConnectionMode) {
if matches!(self.state, SessionState::Idle | SessionState::Scanning) {
self.mode = mode;
}
}
pub fn set_scanning(&mut self) {
if self.state == SessionState::Idle {
self.state = SessionState::Scanning;
}
}
pub fn set_idle_if_scanning(&mut self) {
if self.state == SessionState::Scanning {
self.state = SessionState::Idle;
}
}
pub fn connect_ble(&mut self, device_ref: &str, side: SideKey) {
self.close_session();
self.mode = ConnectionMode::Ble;
self.state = SessionState::Connecting;
self.target = device_ref.into();
self.side = side;
let cancel = CancellationToken::new();
let (incoming_tx, incoming_rx) = mpsc::channel::<Vec<u8>>(128);
let (write_tx, write_rx) = mpsc::channel::<Vec<u8>>(32);
let (event_tx, event_rx) = mpsc::channel::<SessionEvent>(8);
let writer: Arc<Mutex<Box<FrameWriter>>> = Arc::new(Mutex::new(Box::new({
let write_tx = write_tx.clone();
move |frame: &[u8]| write_tx.try_send(frame.to_vec()).is_ok()
})));
self.writer = Some(writer.clone());
self.cancel = Some(cancel.clone());
self.ble_rx = Some(incoming_rx);
self.event_rx = Some(event_rx);
let device_ref = device_ref.to_string();
let ble_event_tx = event_tx.clone();
let ble_cancel = cancel.clone();
spawn_session_thread("BLE", event_tx.clone(), move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build();
if let Ok(rt) = runtime {
let result = rt.block_on(BleTransportAdapter::connect_and_listen(
&device_ref,
incoming_tx,
write_rx,
ble_event_tx.clone(),
ble_cancel,
side,
writer,
));
if let Err(err) = result {
let _ = ble_event_tx
.blocking_send(SessionEvent::Error(format!("BLE connect failed: {err}")));
}
} else {
let _ = ble_event_tx
.blocking_send(SessionEvent::Error("BLE runtime init failed".into()));
}
});
}
pub fn connect_usb(&mut self, port: &str, side: SideKey) {
self.close_session();
self.mode = ConnectionMode::UsbDongle;
self.state = SessionState::Connecting;
self.target = port.into();
self.side = side;
let cancel = CancellationToken::new();
let port = port.to_string();
let (incoming_tx, incoming_rx) = mpsc::channel::<usb_codec::UsbUplinkResult>(128);
let (write_tx, mut write_rx) = mpsc::channel::<Vec<u8>>(32);
let (event_tx, event_rx) = mpsc::channel::<SessionEvent>(8);
let writer: Arc<Mutex<Box<FrameWriter>>> = Arc::new(Mutex::new(Box::new({
let write_tx = write_tx.clone();
move |frame: &[u8]| write_tx.try_send(usb_codec::encode_downlink(frame)).is_ok()
})));
self.writer = Some(writer.clone());
self.usb_rx = Some(incoming_rx);
self.event_rx = Some(event_rx);
self.cancel = Some(cancel.clone());
let hb_cancel = cancel.clone();
spawn_session_thread("USB", event_tx.clone(), move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build();
if let Ok(rt) = runtime {
match rt.block_on(SerialTransportAdapter::open(&port, incoming_tx, cancel)) {
Ok(adapter) => {
let _ = event_tx
.blocking_send(SessionEvent::Ready(format!("{port} 115200 8N1")));
// 串口真正打开后再启动心跳,避免连接建立前写队列无人消费导致误报异常。
HeartbeatService::start(
side,
ConnectionMode::UsbDongle,
writer,
hb_cancel,
Some(event_tx.clone()),
);
while let Some(frame) = write_rx.blocking_recv() {
if let Err(err) = adapter.write_usb_frame(&frame) {
let _ = event_tx.blocking_send(SessionEvent::Error(format!(
"USB write failed: {err}"
)));
break;
}
}
}
Err(err) => {
let _ = event_tx
.blocking_send(SessionEvent::Error(format!("USB open failed: {err}")));
}
}
} else {
let _ =
event_tx.blocking_send(SessionEvent::Error("USB runtime init failed".into()));
}
});
}
pub fn disconnect(&mut self) {
self.close_session();
}
pub fn set_streaming(&mut self, enabled: bool) {
self.state = if enabled {
SessionState::Streaming
} else {
SessionState::Connected
};
}
pub fn set_recording(&mut self, enabled: bool) {
self.state = if enabled {
SessionState::Recording
} else {
SessionState::Connected
};
}
pub fn set_error(&mut self) {
self.state = SessionState::Error;
}
pub fn set_connected(&mut self) {
self.state = SessionState::Connected;
}
pub fn send_frame(&self, frame: &[u8]) -> Result<(), String> {
if let Some(writer) = &self.writer {
let mut guard = writer
.lock()
.map_err(|_| "writer lock poisoned".to_string())?;
if guard(frame) {
Ok(())
} else {
Err("send failed".into())
}
} else {
Err("no active session".into())
}
}
pub fn drain_ble_frames(&mut self) -> Vec<Vec<u8>> {
let mut frames = Vec::new();
if let Some(rx) = self.ble_rx.as_mut() {
while let Ok(frame) = rx.try_recv() {
frames.push(frame);
}
}
frames
}
pub fn drain_usb_frames(&mut self) -> Vec<usb_codec::UsbUplinkResult> {
let mut frames = Vec::new();
if let Some(rx) = self.usb_rx.as_mut() {
while let Ok(frame) = rx.try_recv() {
frames.push(frame);
}
}
frames
}
pub fn drain_events(&mut self) -> Vec<SessionEvent> {
let mut events = Vec::new();
if let Some(rx) = self.event_rx.as_mut() {
while let Ok(event) = rx.try_recv() {
events.push(event);
}
}
events
}
fn close_session(&mut self) {
if let Some(cancel) = self.cancel.take() {
cancel.cancel();
}
self.writer = None;
self.ble_rx = None;
self.usb_rx = None;
self.event_rx = None;
self.state = SessionState::Idle;
}
}
fn spawn_session_thread<F>(label: &'static str, event_tx: mpsc::Sender<SessionEvent>, work: F)
where
F: FnOnce() + Send + 'static,
{
std::thread::spawn(move || {
let result = catch_unwind(AssertUnwindSafe(work));
if let Err(payload) = result {
let _ = event_tx.blocking_send(SessionEvent::Error(format!(
"{label} task panic: {}",
panic_message(payload)
)));
}
});
}
fn panic_message(payload: Box<dyn Any + Send>) -> String {
if let Some(message) = payload.downcast_ref::<&str>() {
(*message).to_string()
} else if let Some(message) = payload.downcast_ref::<String>() {
message.clone()
} else {
"unknown panic".into()
}
}
+220
View File
@@ -0,0 +1,220 @@
use crate::models::{MotorDebugSample, SessionMetrics};
use std::collections::VecDeque;
#[derive(Default)]
pub struct StatisticsService {
window: usize,
samples: VecDeque<StatsPoint>,
total_samples: u64,
seq_gap_count: u64,
seq_out_of_order_count: u64,
protection_state_count: u64,
}
#[derive(Debug, Clone, Copy)]
struct StatsPoint {
host_ts_ms: u64,
seq_id: Option<u16>,
usb_dt_ms: Option<u16>,
}
impl StatisticsService {
pub fn new(window: usize) -> Self {
Self {
window,
samples: VecDeque::new(),
total_samples: 0,
seq_gap_count: 0,
seq_out_of_order_count: 0,
protection_state_count: 0,
}
}
pub fn consume(&mut self, sample: &MotorDebugSample) {
self.total_samples += 1;
if sample.system_state.as_deref().is_some_and(|state| {
let state = state.to_ascii_lowercase();
state.contains("protect")
|| state.contains("fault")
|| state.contains("error")
|| state != "0"
}) {
self.protection_state_count += 1;
}
let seq_id = sample.usb_meta.map(|m| m.seq_id);
if let Some(last) = self.samples.back().and_then(|point| point.seq_id) {
if let Some(curr) = seq_id {
if curr < last {
self.seq_out_of_order_count += 1;
} else {
let expected = last.wrapping_add(1);
if curr != expected {
self.seq_gap_count += 1;
}
}
}
}
self.samples.push_back(StatsPoint {
host_ts_ms: sample.host_ts_ms,
seq_id,
usb_dt_ms: sample.usb_meta.map(|m| m.dt_ms),
});
while self.samples.len() > self.window {
self.samples.pop_front();
}
}
pub fn live_metrics(&self) -> SessionMetrics {
let intervals: Vec<f64> = self
.samples
.iter()
.zip(self.samples.iter().skip(1))
.filter_map(|(prev, curr)| curr.host_ts_ms.checked_sub(prev.host_ts_ms))
.map(|dt| dt as f64)
.collect();
let avg = if intervals.is_empty() {
None
} else {
Some(intervals.iter().sum::<f64>() / intervals.len() as f64)
};
let max = intervals.iter().copied().reduce(f64::max);
let usb_intervals: Vec<f64> = self
.samples
.iter()
.filter_map(|point| point.usb_dt_ms.map(|dt| dt as f64))
.collect();
let usb_avg = if usb_intervals.is_empty() {
None
} else {
Some(usb_intervals.iter().sum::<f64>() / usb_intervals.len() as f64)
};
let usb_max = usb_intervals.iter().copied().reduce(f64::max);
let dropped = if intervals.is_empty() {
None
} else {
let dropped = intervals.iter().filter(|dt| **dt > 30.0).count() as f64;
Some(dropped / intervals.len() as f64)
};
SessionMetrics {
sample_count: self.total_samples,
dropped_frame_rate: dropped,
avg_frame_interval_ms: avg,
max_frame_interval_ms: max,
avg_usb_dt_ms: usb_avg,
max_usb_dt_ms: usb_max,
protection_state_count: self.protection_state_count,
usb_seq_gap_count: self.seq_gap_count,
usb_out_of_order_count: self.seq_out_of_order_count,
}
}
pub fn reset(&mut self) {
self.samples.clear();
self.total_samples = 0;
self.seq_gap_count = 0;
self.seq_out_of_order_count = 0;
self.protection_state_count = 0;
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::models::{ConnectionMode, MotorDebugSample, SideKey};
#[test]
fn counts_usb_seq_gaps() {
let mut stats = StatisticsService::new(10);
let sample = |seq, ts| MotorDebugSample {
host_ts_ms: ts,
connection_mode: ConnectionMode::UsbDongle,
target_id: "x".into(),
side_key: SideKey::Right,
raw_debug_string: String::new(),
current: None,
velocity: None,
position: None,
voltage: None,
motor_temperature: None,
system_state: None,
error_text: None,
usb_meta: Some(crate::models::UsbUplinkMeta {
seq_id: seq,
dt_ms: 10,
dev_id: 0x31,
payload_len: 0,
}),
};
stats.consume(&sample(1, 0));
stats.consume(&sample(3, 50));
let metrics = stats.live_metrics();
assert_eq!(metrics.usb_seq_gap_count, 1);
assert_eq!(metrics.sample_count, 2);
assert_eq!(metrics.max_frame_interval_ms, Some(50.0));
assert_eq!(metrics.dropped_frame_rate, Some(1.0));
}
#[test]
fn counts_usb_out_of_order() {
let mut stats = StatisticsService::new(10);
let sample = |seq, ts| MotorDebugSample {
host_ts_ms: ts,
connection_mode: ConnectionMode::UsbDongle,
target_id: "x".into(),
side_key: SideKey::Right,
raw_debug_string: String::new(),
current: None,
velocity: None,
position: None,
voltage: None,
motor_temperature: None,
system_state: None,
error_text: None,
usb_meta: Some(crate::models::UsbUplinkMeta {
seq_id: seq,
dt_ms: 10,
dev_id: 0x31,
payload_len: 0,
}),
};
stats.consume(&sample(2, 0));
stats.consume(&sample(1, 10));
assert_eq!(stats.live_metrics().usb_out_of_order_count, 1);
}
#[test]
fn reports_usb_dt_separately_from_host_interval() {
let mut stats = StatisticsService::new(10);
let sample = |seq, host_ts, dt_ms| MotorDebugSample {
host_ts_ms: host_ts,
connection_mode: ConnectionMode::UsbDongle,
target_id: "x".into(),
side_key: SideKey::Right,
raw_debug_string: String::new(),
current: None,
velocity: None,
position: None,
voltage: None,
motor_temperature: None,
system_state: None,
error_text: None,
usb_meta: Some(crate::models::UsbUplinkMeta {
seq_id: seq,
dt_ms,
dev_id: 0x31,
payload_len: 0,
}),
};
stats.consume(&sample(1, 100, 20));
stats.consume(&sample(2, 160, 22));
let metrics = stats.live_metrics();
assert_eq!(metrics.avg_frame_interval_ms, Some(60.0));
assert_eq!(metrics.avg_usb_dt_ms, Some(21.0));
assert_eq!(metrics.max_usb_dt_ms, Some(22.0));
}
}
+190
View File
@@ -0,0 +1,190 @@
use crate::models::{DeviceCandidate, SideKey};
use crate::services::heartbeat::{FrameWriter, HeartbeatService};
use crate::services::session::SessionEvent;
use btleplug::api::{Central, CentralEvent, Manager as _, Peripheral as _, ScanFilter, WriteType};
use btleplug::platform::Manager;
use futures::StreamExt;
use std::sync::{Arc, Mutex};
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
pub struct BleTransportAdapter;
impl BleTransportAdapter {
pub async fn start_scan(
device_tx: mpsc::Sender<DeviceCandidate>,
cancel: CancellationToken,
) -> Result<(), String> {
let manager = Manager::new().await.map_err(|err| err.to_string())?;
let adapters = manager.adapters().await.map_err(|err| err.to_string())?;
let adapter = adapters
.into_iter()
.next()
.ok_or_else(|| "No BLE adapter found".to_string())?;
adapter
.start_scan(ScanFilter::default())
.await
.map_err(|err| err.to_string())?;
let mut events = adapter.events().await.map_err(|err| err.to_string())?;
loop {
tokio::select! {
_ = cancel.cancelled() => {
let _ = adapter.stop_scan().await;
break;
}
maybe_event = events.next() => {
let Some(event) = maybe_event else { break };
if let CentralEvent::DeviceDiscovered(id) | CentralEvent::DeviceUpdated(id) = event {
if let Ok(peripheral) = adapter.peripheral(&id).await {
if let Ok(Some(properties)) = peripheral.properties().await {
let name = properties.local_name.unwrap_or_else(|| "(unknown)".into());
let candidate = DeviceCandidate {
device_name: name.clone(),
device_ref: id.to_string(),
rssi: properties.rssi.map(|r| r as i32),
discovered_at_ms: current_time_ms(),
is_preferred: is_preferred_name(&name),
side_hint: extract_side_hint(&name),
};
let _ = device_tx.send(candidate).await;
}
}
}
}
}
}
Ok(())
}
pub async fn connect_and_listen(
device_ref: &str,
frame_tx: mpsc::Sender<Vec<u8>>,
mut write_rx: mpsc::Receiver<Vec<u8>>,
event_tx: mpsc::Sender<SessionEvent>,
cancel: CancellationToken,
side: SideKey,
heartbeat_writer: Arc<Mutex<Box<FrameWriter>>>,
) -> Result<(), String> {
let manager = Manager::new().await.map_err(|err| err.to_string())?;
let adapter = manager
.adapters()
.await
.map_err(|err| err.to_string())?
.into_iter()
.next()
.ok_or_else(|| "No BLE adapter found".to_string())?;
let peripheral = find_peripheral_by_ref(&adapter, device_ref).await?;
peripheral.connect().await.map_err(|err| err.to_string())?;
peripheral
.discover_services()
.await
.map_err(|err| err.to_string())?;
let chars = peripheral.characteristics();
let rx_uuid = Uuid::parse_str("6e400003-b5a3-f393-e0a9-e50e24dcca9e")
.map_err(|err| err.to_string())?;
let tx_uuid = Uuid::parse_str("6e400002-b5a3-f393-e0a9-e50e24dcca9e")
.map_err(|err| err.to_string())?;
let notify_char = chars
.iter()
.find(|c| c.uuid == rx_uuid)
.cloned()
.ok_or_else(|| "RX characteristic missing".to_string())?;
let write_char = chars
.iter()
.find(|c| c.uuid == tx_uuid)
.cloned()
.ok_or_else(|| "TX characteristic missing".to_string())?;
peripheral
.subscribe(&notify_char)
.await
.map_err(|err| err.to_string())?;
let mut notifications = peripheral
.notifications()
.await
.map_err(|err| err.to_string())?;
let _ = event_tx
.send(SessionEvent::Ready(format!("{device_ref} BLE NUS")))
.await;
HeartbeatService::start(
side,
crate::models::ConnectionMode::Ble,
heartbeat_writer,
cancel.clone(),
Some(event_tx.clone()),
);
loop {
tokio::select! {
_ = cancel.cancelled() => {
let _ = peripheral.disconnect().await;
break;
}
maybe_note = notifications.next() => {
let Some(note) = maybe_note else { break };
let _ = frame_tx.send(note.value).await;
}
maybe_frame = write_rx.recv() => {
let Some(frame) = maybe_frame else { break };
peripheral
.write(&write_char, &frame, WriteType::WithoutResponse)
.await
.map_err(|err| err.to_string())?;
}
}
}
let _ = peripheral.unsubscribe(&notify_char).await;
let _ = peripheral.disconnect().await;
let _ = write_char;
Ok(())
}
}
pub fn extract_side_hint(name: &str) -> Option<SideKey> {
let upper = name.to_uppercase();
if upper.contains("LEFT")
|| upper.contains("ZDL")
|| upper.contains("ZCL")
|| upper.contains("ARL")
{
Some(SideKey::Left)
} else if upper.contains("RIGHT")
|| upper.contains("ZDR")
|| upper.contains("ZCR")
|| upper.contains("ZDB")
{
Some(SideKey::Right)
} else {
None
}
}
fn is_preferred_name(name: &str) -> bool {
let upper = name.to_uppercase();
upper.contains("YROBOT")
}
fn current_time_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
async fn find_peripheral_by_ref(
adapter: &btleplug::platform::Adapter,
device_ref: &str,
) -> Result<btleplug::platform::Peripheral, String> {
for peripheral in adapter.peripherals().await.map_err(|err| err.to_string())? {
if peripheral.id().to_string() == device_ref
|| peripheral.address().to_string() == device_ref
{
return Ok(peripheral);
}
}
Err(format!("BLE device not found: {device_ref}"))
}
+2
View File
@@ -0,0 +1,2 @@
pub mod ble_adapter;
pub mod serial_adapter;
+73
View File
@@ -0,0 +1,73 @@
use crate::protocol::usb_codec::{UsbStreamAssembler, UsbUplinkResult};
use std::io::{Read, Write};
use std::sync::{Arc, Mutex};
use std::thread;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
pub struct SerialTransportAdapter {
port: Arc<Mutex<Box<dyn serialport::SerialPort>>>,
}
impl SerialTransportAdapter {
pub async fn open(
port: &str,
result_tx: mpsc::Sender<UsbUplinkResult>,
cancel: CancellationToken,
) -> Result<Box<Self>, String> {
let serial = serialport::new(port, 115_200)
.data_bits(serialport::DataBits::Eight)
.flow_control(serialport::FlowControl::None)
.parity(serialport::Parity::None)
.stop_bits(serialport::StopBits::One)
.timeout(std::time::Duration::from_millis(100))
.open()
.map_err(|err| err.to_string())?;
let adapter = Box::new(Self {
port: Arc::new(Mutex::new(serial)),
});
let reader_port = adapter.port.clone();
thread::spawn(move || {
let mut assembler = UsbStreamAssembler::default();
let mut buf = [0u8; 512];
while !cancel.is_cancelled() {
let read_result = {
let mut guard = match reader_port.lock() {
Ok(guard) => guard,
Err(_) => break,
};
guard.read(&mut buf)
};
match read_result {
Ok(n) if n > 0 => {
for result in assembler.push(&buf[..n]) {
let _ = result_tx.blocking_send(result);
}
}
Ok(_) => {}
Err(err) if err.kind() == std::io::ErrorKind::TimedOut => {}
Err(err) => {
let _ = result_tx.blocking_send(UsbUplinkResult {
meta: None,
payload: vec![],
error: Some(err.to_string()),
});
break;
}
}
}
});
Ok(adapter)
}
pub fn write_usb_frame(&self, frame: &[u8]) -> Result<(), String> {
let mut guard = self
.port
.lock()
.map_err(|_| "Serial lock poisoned".to_string())?;
guard.write_all(frame).map_err(|err| err.to_string())
}
}
+20
View File
@@ -0,0 +1,20 @@
use crate::models::{ConnectionMode, SideKey};
#[derive(Debug, Clone)]
pub enum UiAction {
SelectMode(ConnectionMode),
Scan,
StopScan,
RefreshPorts,
SelectSide(SideKey),
ConnectBle(String),
ConnectUsb(String),
Disconnect,
SetParams(serde_json::Value),
QueryCurrent,
QueryMode(u8),
EnableDebug,
DisableDebug,
StartRecording(String),
StopRecording,
}
+131
View File
@@ -0,0 +1,131 @@
use crate::models::{ConnectionMode, DeviceCandidate, SessionState, SideKey, UiSessionViewState};
use crate::ui::actions::UiAction;
use eframe::egui;
use eframe::egui::Ui;
#[derive(Default)]
pub struct ConnectionPanel {
pub mode: ConnectionMode,
scan_items: Vec<DeviceCandidate>,
ports: Vec<String>,
usb_port: String,
selected_ble_ref: String,
side: SideKey,
session_state: SessionState,
}
impl ConnectionPanel {
pub fn set_connection_state(&mut self, state: &UiSessionViewState) {
self.mode = state.connection_mode;
self.session_state = state.session_state;
if let Some(side) = state.side_key {
self.side = side;
}
}
pub fn ui(&mut self, ui: &mut Ui, actions: &mut Vec<UiAction>) {
ui.horizontal(|ui| {
ui.label("Mode:");
egui::ComboBox::from_id_salt("mode_selector")
.selected_text(match self.mode {
ConnectionMode::Ble => "BLE",
ConnectionMode::UsbDongle => "USB Dongle",
})
.show_ui(ui, |ui| {
if ui
.selectable_value(&mut self.mode, ConnectionMode::Ble, "BLE")
.clicked()
{
actions.push(UiAction::SelectMode(ConnectionMode::Ble));
}
if ui
.selectable_value(&mut self.mode, ConnectionMode::UsbDongle, "USB Dongle")
.clicked()
{
actions.push(UiAction::SelectMode(ConnectionMode::UsbDongle));
}
});
});
ui.horizontal(|ui| {
ui.label("Side:");
for (side, label) in [
(SideKey::Left, "Left"),
(SideKey::Right, "Right"),
(SideKey::Neutral, "Neutral"),
] {
if ui.selectable_value(&mut self.side, side, label).clicked() {
actions.push(UiAction::SelectSide(side));
}
}
});
let can_connect = matches!(
self.session_state,
SessionState::Idle | SessionState::Scanning
);
match self.mode {
ConnectionMode::Ble => {
ui.horizontal(|ui| {
if ui.button("Scan / Refresh").clicked() {
self.scan_items.clear();
actions.push(UiAction::Scan);
}
if ui.button("Stop Scan").clicked() {
actions.push(UiAction::StopScan);
}
if ui
.add_enabled(can_connect, egui::Button::new("Connect BLE"))
.clicked()
{
actions.push(UiAction::ConnectBle(self.selected_ble_ref.clone()));
}
});
}
ConnectionMode::UsbDongle => {
ui.horizontal(|ui| {
ui.label("USB Port:");
ui.text_edit_singleline(&mut self.usb_port);
if ui.button("Refresh Ports").clicked() {
actions.push(UiAction::RefreshPorts);
}
if ui
.add_enabled(can_connect, egui::Button::new("Open USB"))
.clicked()
{
actions.push(UiAction::ConnectUsb(self.usb_port.clone()));
}
});
}
}
if ui.button("Disconnect").clicked() {
actions.push(UiAction::Disconnect);
}
for item in &self.scan_items {
let label = format!(
"{} {} RSSI:{:?}{}",
item.device_name,
item.device_ref,
item.rssi,
if item.is_preferred { " preferred" } else { "" }
);
if ui
.selectable_label(self.selected_ble_ref == item.device_ref, label)
.clicked()
{
self.selected_ble_ref = item.device_ref.clone();
}
}
for port in &self.ports {
if ui.selectable_label(self.usb_port == *port, port).clicked() {
self.usb_port = port.clone();
}
}
}
pub fn set_scan_items(&mut self, items: Vec<DeviceCandidate>) {
self.scan_items = items;
}
pub fn set_ports(&mut self, ports: Vec<String>) {
self.ports = ports;
}
}
+182
View File
@@ -0,0 +1,182 @@
use crate::models::{SessionState, UiSessionViewState};
use crate::ui::actions::UiAction;
use eframe::egui;
use eframe::egui::Ui;
use serde_json::json;
pub struct ControlPanel {
mode_index: usize,
fm_mset: f64,
fm_vset: f64,
fm_pset: f64,
fm_kp: f64,
fm_ki: f64,
fm_kd: f64,
session_state: SessionState,
}
impl Default for ControlPanel {
fn default() -> Self {
Self {
mode_index: 0,
fm_mset: 0.0,
fm_vset: 0.0,
fm_pset: 0.0,
fm_kp: 0.0,
fm_ki: 0.0,
fm_kd: 0.0,
session_state: SessionState::Idle,
}
}
}
impl ControlPanel {
pub fn set_session_state(&mut self, state: &UiSessionViewState) {
self.session_state = state.session_state;
}
pub fn ui(&mut self, ui: &mut Ui, actions: &mut Vec<UiAction>) {
ui.horizontal(|ui| {
ui.label("Force mode:");
egui::ComboBox::from_id_salt("force_mode")
.selected_text(self.mode_index.to_string())
.show_ui(ui, |ui| {
for mode in 0..=6 {
ui.selectable_value(&mut self.mode_index, mode, mode.to_string());
}
});
});
self.params_grid(ui);
let connected = matches!(
self.session_state,
SessionState::Connected | SessionState::Streaming | SessionState::Recording
);
ui.horizontal(|ui| {
if ui
.add_enabled(connected, egui::Button::new("Set Params"))
.clicked()
{
actions.push(UiAction::SetParams(self.params_json()));
}
if ui
.add_enabled(connected, egui::Button::new("Query Current"))
.clicked()
{
actions.push(UiAction::QueryCurrent);
}
if ui
.add_enabled(connected, egui::Button::new("Query Mode"))
.clicked()
{
actions.push(UiAction::QueryMode(self.mode_index as u8));
}
if ui
.add_enabled(connected, egui::Button::new("Enable Realtime"))
.clicked()
{
actions.push(UiAction::EnableDebug);
}
if ui
.add_enabled(connected, egui::Button::new("Disable Realtime"))
.clicked()
{
actions.push(UiAction::DisableDebug);
}
});
}
fn params_grid(&mut self, ui: &mut Ui) {
egui::Grid::new("force_params_grid")
.num_columns(3)
.striped(true)
.show(ui, |ui| {
// 按模式只暴露相关字段,避免把上一模式的旧参数误发给固件。
for field in Self::fields_for_mode(self.mode_index as u8) {
ui.label(field.label);
ui.add(
egui::DragValue::new(field.value_mut(self))
.speed(field.speed)
.range(0.0..=f64::MAX),
);
ui.label(field.unit);
ui.end_row();
}
});
}
fn params_json(&self) -> serde_json::Value {
match self.mode_index {
0 => json!({ "fm_mode": 0 }),
1 => json!({ "fm_mode": 1, "fm_mset": self.fm_mset }),
2 => json!({ "fm_mode": 2, "fm_vset": self.fm_vset }),
3 => json!({ "fm_mode": 3, "fm_pset": self.fm_pset }),
4 => json!({ "fm_mode": 4, "fm_mset": self.fm_mset, "fm_vset": self.fm_vset }),
5 => {
json!({ "fm_mode": 5, "fm_mset": self.fm_mset, "fm_kp": self.fm_kp, "fm_kd": self.fm_kd })
}
_ => json!({
"fm_mode": 6,
"fm_mset": self.fm_mset,
"fm_vset": self.fm_vset,
"fm_pset": self.fm_pset,
"fm_kp": self.fm_kp,
"fm_ki": self.fm_ki,
"fm_kd": self.fm_kd
}),
}
}
fn fields_for_mode(mode: u8) -> Vec<ForceParamField> {
match mode {
0 => vec![],
1 => vec![ForceParamField::new("fm_mset", "N")],
2 => vec![ForceParamField::new("fm_vset", "rad/s")],
3 => vec![ForceParamField::new("fm_pset", "rad")],
4 => vec![
ForceParamField::new("fm_mset", "N"),
ForceParamField::new("fm_vset", "rad/s"),
],
5 => vec![
ForceParamField::new("fm_mset", "N"),
ForceParamField::new("fm_kp", ""),
ForceParamField::new("fm_kd", ""),
],
_ => vec![
ForceParamField::new("fm_mset", "N"),
ForceParamField::new("fm_vset", "rad/s"),
ForceParamField::new("fm_pset", "rad"),
ForceParamField::new("fm_kp", ""),
ForceParamField::new("fm_ki", ""),
ForceParamField::new("fm_kd", ""),
],
}
}
}
struct ForceParamField {
label: &'static str,
unit: &'static str,
speed: f64,
}
impl ForceParamField {
fn new(label: &'static str, unit: &'static str) -> Self {
Self {
label,
unit,
speed: 0.1,
}
}
fn value_mut<'a>(&self, panel: &'a mut ControlPanel) -> &'a mut f64 {
match self.label {
"fm_mset" => &mut panel.fm_mset,
"fm_vset" => &mut panel.fm_vset,
"fm_pset" => &mut panel.fm_pset,
"fm_kp" => &mut panel.fm_kp,
"fm_ki" => &mut panel.fm_ki,
"fm_kd" => &mut panel.fm_kd,
_ => &mut panel.fm_mset,
}
}
}
+45
View File
@@ -0,0 +1,45 @@
use eframe::egui::{ScrollArea, Ui};
#[derive(Default)]
pub struct EventLogPanel {
events: Vec<String>,
}
impl EventLogPanel {
pub fn append_event(&mut self, msg: &str) {
self.events.push(msg.to_string());
if self.events.len() > 500 {
self.events.remove(0);
}
}
#[cfg(test)]
pub fn events(&self) -> &[String] {
&self.events
}
pub fn ui(&mut self, ui: &mut Ui) {
ScrollArea::vertical().max_height(100.0).show(ui, |ui| {
for evt in &self.events {
ui.label(evt);
}
});
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn keeps_recent_500_events() {
let mut panel = EventLogPanel::default();
for idx in 0..510 {
panel.append_event(&format!("event-{idx}"));
}
assert_eq!(panel.events().len(), 500);
assert_eq!(panel.events()[0], "event-10");
assert_eq!(panel.events()[499], "event-509");
}
}
+108
View File
@@ -0,0 +1,108 @@
use eframe::egui;
use egui::{CentralPanel, TopBottomPanel};
use crate::models::{DeviceCandidate, UiSessionViewState};
use crate::ui::actions::UiAction;
use crate::ui::connection_panel::ConnectionPanel;
use crate::ui::control_panel::ControlPanel;
use crate::ui::event_log_panel::EventLogPanel;
use crate::ui::realtime_panel::RealtimePanel;
use crate::ui::recording_panel::RecordingPanel;
use crate::ui::status_panel::StatusPanel;
pub struct MainWindow {
pub connection_panel: ConnectionPanel,
pub control_panel: ControlPanel,
pub status_panel: StatusPanel,
pub realtime_panel: RealtimePanel,
pub recording_panel: RecordingPanel,
pub event_log_panel: EventLogPanel,
pub status_label: String,
pub view_state: UiSessionViewState,
pub error_msg: String,
}
impl Default for MainWindow {
fn default() -> Self {
Self {
connection_panel: ConnectionPanel::default(),
control_panel: ControlPanel::default(),
status_panel: StatusPanel::default(),
realtime_panel: RealtimePanel::default(),
recording_panel: RecordingPanel::default(),
event_log_panel: EventLogPanel::default(),
status_label: "Ready".into(),
view_state: UiSessionViewState::default(),
error_msg: String::new(),
}
}
}
impl MainWindow {
pub fn update_view_state(&mut self, state: UiSessionViewState) {
self.view_state = state.clone();
self.status_label = format!("{:?}", state.session_state);
self.connection_panel.set_connection_state(&state);
self.control_panel.set_session_state(&state);
self.status_panel.update_status(&state);
self.recording_panel.set_recording_state(&state);
}
pub fn show_error(&mut self, msg: &str) {
self.error_msg = msg.to_string();
self.status_label = format!("ERROR: {msg}");
self.event_log_panel.append_event(&format!("ERROR: {msg}"));
}
pub fn set_scan_items(&mut self, items: Vec<DeviceCandidate>) {
self.connection_panel.set_scan_items(items);
}
pub fn set_ports(&mut self, ports: Vec<String>) {
self.connection_panel.set_ports(ports);
}
pub fn ui(&mut self, ctx: &egui::Context) -> Vec<UiAction> {
let mut actions = Vec::new();
TopBottomPanel::top("status_bar").show(ctx, |ui| {
ui.horizontal(|ui| {
ui.label(&self.status_label);
ui.separator();
ui.label(format!("{:?}", self.view_state.session_state));
ui.separator();
ui.label(if self.view_state.heartbeat_ok {
"HB OK"
} else {
"HB ERR"
});
ui.separator();
ui.label(if self.view_state.recording_active {
"REC ON"
} else {
"REC OFF"
});
});
});
CentralPanel::default().show(ctx, |ui| {
ui.strong("Connection");
self.connection_panel.ui(ui, &mut actions);
ui.separator();
ui.strong("Motor Control");
self.control_panel.ui(ui, &mut actions);
ui.separator();
ui.strong("Realtime Data");
self.realtime_panel.ui(ui);
ui.separator();
ui.strong("Recording");
self.recording_panel.ui(ui, &mut actions);
ui.separator();
ui.strong("Status");
self.status_panel.ui(ui);
ui.separator();
ui.strong("Event Log");
self.event_log_panel.ui(ui);
});
actions
}
}
+8
View File
@@ -0,0 +1,8 @@
pub mod actions;
pub mod connection_panel;
pub mod control_panel;
pub mod event_log_panel;
pub mod main_window;
pub mod realtime_panel;
pub mod recording_panel;
pub mod status_panel;
+89
View File
@@ -0,0 +1,89 @@
use crate::models::MotorDebugSample;
use eframe::egui::Ui;
use egui_plot::{Line, Plot, PlotPoints};
use std::collections::VecDeque;
pub struct RealtimePanel {
samples: VecDeque<MotorDebugSample>,
show_current: bool,
show_velocity: bool,
show_position: bool,
show_voltage: bool,
show_temperature: bool,
max_samples: usize,
}
impl Default for RealtimePanel {
fn default() -> Self {
Self {
samples: VecDeque::new(),
show_current: true,
show_velocity: true,
show_position: false,
show_voltage: true,
show_temperature: true,
max_samples: 600,
}
}
}
impl RealtimePanel {
pub fn push_sample(&mut self, sample: MotorDebugSample) {
self.samples.push_back(sample);
while self.samples.len() > self.max_samples {
self.samples.pop_front();
}
}
#[cfg(test)]
pub fn sample_count(&self) -> usize {
self.samples.len()
}
pub fn ui(&mut self, ui: &mut Ui) {
ui.horizontal(|ui| {
ui.checkbox(&mut self.show_current, "current");
ui.checkbox(&mut self.show_velocity, "velocity");
ui.checkbox(&mut self.show_position, "position");
ui.checkbox(&mut self.show_voltage, "voltage");
ui.checkbox(&mut self.show_temperature, "temperature");
});
ui.label(format!("Realtime samples: {}", self.samples.len()));
Plot::new("realtime_plot")
.height(220.0)
.show(ui, |plot_ui| {
if self.show_current {
plot_ui.line(line_for(&self.samples, "current", |s| s.current));
}
if self.show_velocity {
plot_ui.line(line_for(&self.samples, "velocity", |s| s.velocity));
}
if self.show_position {
plot_ui.line(line_for(&self.samples, "position", |s| s.position));
}
if self.show_voltage {
plot_ui.line(line_for(&self.samples, "voltage", |s| s.voltage));
}
if self.show_temperature {
plot_ui.line(line_for(&self.samples, "temperature", |s| {
s.motor_temperature
}));
}
});
}
}
fn line_for(
samples: &VecDeque<MotorDebugSample>,
name: &'static str,
value: impl Fn(&MotorDebugSample) -> Option<f64>,
) -> Line<'static> {
let points: PlotPoints<'static> = samples
.iter()
.enumerate()
.filter_map(|(idx, sample)| value(sample).map(|v| [idx as f64, v]))
.collect::<Vec<_>>()
.into();
Line::new(points).name(name)
}
+35
View File
@@ -0,0 +1,35 @@
use crate::models::UiSessionViewState;
use crate::ui::actions::UiAction;
use eframe::egui::Ui;
#[derive(Default)]
pub struct RecordingPanel {
pub csv_path: String,
recording: bool,
}
impl RecordingPanel {
pub fn set_recording_state(&mut self, state: &UiSessionViewState) {
self.recording = state.recording_active;
}
pub fn ui(&mut self, ui: &mut Ui, actions: &mut Vec<UiAction>) {
ui.horizontal(|ui| {
ui.label("CSV Path:");
ui.text_edit_singleline(&mut self.csv_path);
});
ui.horizontal(|ui| {
if ui.button("Start Recording").clicked() {
actions.push(UiAction::StartRecording(self.csv_path.clone()));
}
if ui.button("Stop Recording").clicked() {
actions.push(UiAction::StopRecording);
}
});
ui.label(if self.recording {
"Recording: ON"
} else {
"Recording: OFF"
});
}
}
+55
View File
@@ -0,0 +1,55 @@
use crate::models::UiSessionViewState;
use eframe::egui::Ui;
#[derive(Default)]
pub struct StatusPanel {
state: UiSessionViewState,
}
impl StatusPanel {
pub fn update_status(&mut self, state: &UiSessionViewState) {
self.state = state.clone();
}
pub fn ui(&mut self, ui: &mut Ui) {
ui.label(format!(
"Mode: {:?} State: {:?} Target: {} Side: {:?}",
self.state.connection_mode,
self.state.session_state,
self.state.target_label,
self.state.side_key
));
ui.label(format!(
"Samples: {} USB seq gaps: {}",
self.state.metrics.sample_count, self.state.metrics.usb_seq_gap_count
));
ui.label(format!(
"Host interval avg/max: {}/{} ms USB dt avg/max: {}/{} ms",
fmt_opt(self.state.metrics.avg_frame_interval_ms),
fmt_opt(self.state.metrics.max_frame_interval_ms),
fmt_opt(self.state.metrics.avg_usb_dt_ms),
fmt_opt(self.state.metrics.max_usb_dt_ms),
));
ui.label(format!(
"Drop rate: {} USB out-of-order: {} Protection states: {}",
fmt_percent(self.state.metrics.dropped_frame_rate),
self.state.metrics.usb_out_of_order_count,
self.state.metrics.protection_state_count
));
if let Some(error) = &self.state.latest_error {
ui.label(format!("Latest error: {error}"));
}
}
}
fn fmt_opt(value: Option<f64>) -> String {
value
.map(|v| format!("{v:.1}"))
.unwrap_or_else(|| "-".into())
}
fn fmt_percent(value: Option<f64>) -> String {
value
.map(|v| format!("{:.1}%", v * 100.0))
.unwrap_or_else(|| "-".into())
}