From 1f95b485a7766e627bcd35df7c5b2e73fac6083d Mon Sep 17 00:00:00 2001 From: Sachymetsu Date: Tue, 13 Jan 2026 07:43:40 +0100 Subject: [PATCH] Prototype RPC server impl --- Cargo.lock | 8 ++--- src/detector/data.rs | 10 ++++--- src/main.rs | 1 + src/net.rs | 69 +++++++++++++++++++++++++++++--------------- src/rpc.rs | 30 +++++++++++++++++++ src/updates.rs | 8 ++--- 6 files changed, 90 insertions(+), 36 deletions(-) create mode 100644 src/rpc.rs diff --git a/Cargo.lock b/Cargo.lock index 81ecace..17e6141 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -570,7 +570,7 @@ dependencies = [ "nb 1.1.0", "pio", "rand_core 0.6.4", - "rand_core 0.9.3", + "rand_core 0.9.4", "rp-pac", "rp2040-boot2", "sha2-const-stable", @@ -1402,9 +1402,9 @@ checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" [[package]] name = "rand_core" -version = "0.9.3" +version = "0.9.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "99d9a13982dcf210057a8a78572b2217b667c3beacbf3a0d8b454f6f82837d38" +checksum = "4f1b3bc831f92381018fd9c6350b917c7b21f1eed35a65a51900e0e55a3d7afa" [[package]] name = "redox_syscall" @@ -1693,7 +1693,7 @@ dependencies = [ [[package]] name = "striker-proto" version = "0.1.0" -source = "git+https://tangled.org/sachy.dev/striker#674322b5d31c9f4d2c2e129c4c4f17d72cee2b2d" +source = "git+https://tangled.org/sachy.dev/striker#66ac11c20cfedbe0da31e0ca83ce460143e5aa74" dependencies = [ "postcard", "serde", diff --git a/src/detector/data.rs b/src/detector/data.rs index 386fb92..fe3a31e 100644 --- a/src/detector/data.rs +++ b/src/detector/data.rs @@ -1,11 +1,13 @@ +use striker_proto::{StrikerResponse, Update}; + use crate::updates::NetDataSender; pub(super) fn transmit_level_update(timestamp: i64, warn_level: u16, net_data: &NetDataSender) { net_data - .try_send(striker_proto::Update::Warning { + .try_send(StrikerResponse::Update(Update::Warning { timestamp, level: warn_level, - }) + })) .ok(); } @@ -17,7 +19,7 @@ pub(super) fn transmit_strike( net_data: &NetDataSender, ) { net_data - .try_send(striker_proto::Update::Strike { + .try_send(StrikerResponse::Update(Update::Strike { timestamp, peaks: peaks.to_vec(), samples: samples @@ -25,6 +27,6 @@ pub(super) fn transmit_strike( .map(|&sample| (average as i16).saturating_sub_unsigned(sample)) .collect(), average, - }) + })) .ok(); } diff --git a/src/main.rs b/src/main.rs index d06e1cf..7e6c68a 100644 --- a/src/main.rs +++ b/src/main.rs @@ -13,6 +13,7 @@ mod net; #[cfg(not(feature = "defmt"))] mod panic_handler; mod pwm; +mod rpc; mod rtc; mod state; mod updates; diff --git a/src/net.rs b/src/net.rs index fd4fead..3d9ec74 100644 --- a/src/net.rs +++ b/src/net.rs @@ -2,7 +2,7 @@ use alloc::vec; use embassy_futures::select::{select, select3}; use embassy_net::{ dns, - tcp::TcpSocket, + tcp::{TcpReader, TcpSocket, TcpWriter}, udp::{PacketMetadata, UdpSocket}, }; use embassy_time::{Duration, Timer, WithTimeout}; @@ -16,6 +16,7 @@ use sachy_sntp::SntpSocket; use crate::{ constants::{HOST_NAME, HOST_PORT}, + rpc::RpcServer, rtc::GlobalRtc, updates::{NetDataReceiver, UpdateConnection}, utils::try_static_buffer_with, @@ -129,11 +130,13 @@ pub async fn tcp_stack(stack: embassy_net::Stack<'static>) { select(stack.wait_link_down(), data_loop(&mut tcp, &net_data)).await; + net_data.clear(); + UpdateConnection::disconnect(); } } -async fn data_loop<'connection>(tcp: &mut TcpSocket<'connection>, net_data: &NetDataReceiver) { +async fn data_loop<'device>(tcp: &mut TcpSocket<'device>, net_data: &NetDataReceiver) { loop { UpdateConnection::disconnect(); @@ -144,36 +147,54 @@ async fn data_loop<'connection>(tcp: &mut TcpSocket<'connection>, net_data: &Net info!("Connected!"); UpdateConnection::connect(); - 'inner: loop { - let data = net_data.receive().await; + let (reader, writer) = tcp.split(); - if !tcp.may_send() { - // Clear backlog, no point in keeping the updates - // if there is nothing to send the updates to - net_data.clear(); - break 'inner; - } + select(read_loop(reader), write_loop(writer, net_data)).await; + + net_data.clear(); + + info!("DISCONNECT"); + } +} - if tcp - .write_with(|buf| { - let written = unwrap!(postcard::to_slice( - &striker_proto::StrikerResponse::Update(data), - buf - )); +async fn read_loop<'device>(mut reader: TcpReader<'device>) { + let mut buf = vec![0u8; 2048]; - (written.len(), ()) - }) + loop { + if let Ok(read) = reader.read(&mut buf).await + && let Ok(req) = striker_proto::receive_request(&mut buf[..read]) + { + match RpcServer::handle_request(req) .await - .is_err() + .zip(UpdateConnection::can_update()) { - break 'inner; + Some((resp, resp_tx)) => { + resp_tx.try_send(resp).ok(); + } + None => break, } + } + } +} - if tcp.flush().await.is_err() { - break 'inner; - } +async fn write_loop<'device>(mut writer: TcpWriter<'device>, net_data: &NetDataReceiver) { + loop { + let data = net_data.receive().await; + + if writer + .write_with(|buf| { + let written = unwrap!(striker_proto::send_response(data, buf)); + + (written.len(), ()) + }) + .await + .is_err() + { + break; } - info!("DISCONNECT"); + if writer.flush().await.is_err() { + break; + } } } diff --git a/src/rpc.rs b/src/rpc.rs new file mode 100644 index 0000000..e463cc4 --- /dev/null +++ b/src/rpc.rs @@ -0,0 +1,30 @@ +use alloc::string::ToString; +use striker_proto::{Request, Response, StrikerRequest, StrikerResponse}; + +use crate::constants::{BLIP_SIZE, BLIP_THRESHOLD}; + +pub struct RpcServer; + +impl RpcServer { + pub async fn handle_request(req: StrikerRequest) -> Option { + match req.request { + Request::Ping => { + Some(StrikerResponse::Response(Response::Pong)) + } + Request::DetectorInfo => Some(StrikerResponse::Response( + Response::DetectorInfo { + blip_threshold: BLIP_THRESHOLD as usize, + blip_size: BLIP_SIZE, + max_duty: 100, + duty: 0, + }, + )), + Request::SetDetectorConfig { .. } => Some(StrikerResponse::Response( + Response::SetDetectorConfig { + success: false, + message: Some("NOT IMPLEMENTED".to_string()), + }, + )), + } + } +} diff --git a/src/updates.rs b/src/updates.rs index e45e81c..a947917 100644 --- a/src/updates.rs +++ b/src/updates.rs @@ -1,14 +1,14 @@ use embassy_sync::channel::{Channel, Receiver, Sender}; -use striker_proto::Update; +use striker_proto::{StrikerResponse}; use crate::{ locks::NetDataLock, state::{DEVICE_STATE, DeviceState}, }; -pub type NetDataChannel = Channel; -pub type NetDataSender = Sender<'static, NetDataLock, Update, 8>; -pub type NetDataReceiver = Receiver<'static, NetDataLock, Update, 8>; +pub type NetDataChannel = Channel; +pub type NetDataSender = Sender<'static, NetDataLock, StrikerResponse, 8>; +pub type NetDataReceiver = Receiver<'static, NetDataLock, StrikerResponse, 8>; static NET_CHANNEL: NetDataChannel = Channel::new(); -- 2.51.2