diff --git a/Cargo.lock b/Cargo.lock index 17e6141..5acd7c5 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.4", + "rand_core 0.9.5", "rp-pac", "rp2040-boot2", "sha2-const-stable", @@ -1402,9 +1402,9 @@ checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" [[package]] name = "rand_core" -version = "0.9.4" +version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4f1b3bc831f92381018fd9c6350b917c7b21f1eed35a65a51900e0e55a3d7afa" +checksum = "76afc826de14238e6e8c374ddcc1fa19e374fd8dd986b0d2af0d02377261d83c" [[package]] name = "redox_syscall" diff --git a/src/net.rs b/src/net.rs index dad553f..326b37a 100644 --- a/src/net.rs +++ b/src/net.rs @@ -1,24 +1,23 @@ -use alloc::vec; +mod mdns; +mod rpc; +mod sntp; + use embassy_futures::select::{select, select3}; use embassy_net::{ - dns, - tcp::{TcpReader, TcpSocket, TcpWriter}, + tcp::TcpSocket, udp::{PacketMetadata, UdpSocket}, }; -use embassy_time::{Duration, Timer, WithTimeout}; -use sachy_fmt::{error, info, unwrap}; +use embassy_time::Duration; +use sachy_fmt::unwrap; use sachy_mdns::{ - GROUP_ADDR_V4, GROUP_SOCK_V4, MDNS_PORT, + GROUP_ADDR_V4, service::{MdnsService, Service}, - state::MdnsAction, }; -use sachy_sntp::SntpSocket; use crate::{ constants::{HOST_NAME, HOST_PORT}, - rpc::RpcServer, rtc::GlobalRtc, - updates::{NetDataReceiver, UpdateConnection}, + updates::UpdateConnection, utils::try_static_buffer_with, }; @@ -49,12 +48,12 @@ pub async fn udp_stack(stack: embassy_net::Stack<'static>, rtc: GlobalRtc<'stati stack.wait_config_up().await; if !rtc.is_ready().await { - sntp_loop(&mut udp, stack, rtc).await; + sntp::sntp_loop(&mut udp, stack, rtc).await; } select3( stack.wait_link_down(), - mdns_loop(&mut service, &mut udp), + mdns::mdns_loop(&mut service, &mut udp), rtc.wait_for_reset(), ) .await; @@ -63,58 +62,6 @@ pub async fn udp_stack(stack: embassy_net::Stack<'static>, rtc: GlobalRtc<'stati } } -async fn sntp_loop<'device>( - udp: &mut UdpSocket<'device>, - stack: embassy_net::Stack<'device>, - rtc: GlobalRtc<'static>, -) { - loop { - let Ok(addr) = stack.dns_query("pool.ntp.org", dns::DnsQueryType::A).await else { - error!("Failed to query DNS for an NTP server. Retrying..."); - continue; - }; - - match udp.resolve_time(addr.as_slice()).await { - Ok(time) => { - unwrap!(rtc.set_rtc_datetime(time).await); - break; - } - Err(e) => error!("Failed to resolve SNTP time: {}", e), - } - } -} - -async fn mdns_loop<'device>(service: &mut MdnsService, udp: &mut UdpSocket<'device>) { - unwrap!(udp.bind(MDNS_PORT)); - - let mut send_buf = vec![0u8; 2048]; - - loop { - let ready = match service.next_action() { - MdnsAction::Announce => service.send_announcement(&mut send_buf), - - MdnsAction::ListenFor { timeout } => udp - .recv_from_with(|buf, _from| service.listen_for_queries(buf, &mut send_buf)) - .with_timeout(timeout) - .await - .ok() - .flatten(), - - MdnsAction::WaitFor { duration } => { - Timer::after(duration).await; - - None - } - }; - - if let Some(bytes) = ready - && udp.send_to(bytes, GROUP_SOCK_V4).await.is_ok() - { - udp.flush().await; - } - } -} - #[embassy_executor::task] pub async fn tcp_stack(stack: embassy_net::Stack<'static>) { let rx_buffer = unwrap!(try_static_buffer_with(4096, Default::default)); @@ -128,76 +75,10 @@ pub async fn tcp_stack(stack: embassy_net::Stack<'static>) { loop { stack.wait_config_up().await; - select(stack.wait_link_down(), data_loop(&mut tcp, &net_data)).await; + select(stack.wait_link_down(), rpc::data_loop(&mut tcp, &net_data)).await; net_data.clear(); UpdateConnection::disconnect(); } } - -async fn data_loop<'device>(tcp: &mut TcpSocket<'device>, net_data: &NetDataReceiver) { - loop { - UpdateConnection::disconnect(); - - if tcp.accept(HOST_PORT).await.is_err() { - continue; - } - - info!("Connected!"); - UpdateConnection::connect(); - - let (reader, writer) = tcp.split(); - - select(read_loop(reader), write_loop(writer, net_data)).await; - - net_data.clear(); - - info!("DISCONNECT"); - } -} - -async fn read_loop<'device>(mut reader: TcpReader<'device>) { - let mut buf = vec![0u8; 2048]; - - loop { - match reader.read(&mut buf).await { - Ok(0) | Err(_) => break, - Ok(read) => { - if let Ok(req) = striker_proto::receive_request(&mut buf[..read]) - .inspect_err(|e| error!("Proto Error: {}", e)) - { - let Some((resp, resp_tx)) = RpcServer::handle_request(req) - .await - .zip(UpdateConnection::can_update()) - else { - break; - }; - resp_tx.try_send(resp).ok(); - } - } - } - } -} - -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; - } - - if writer.flush().await.is_err() { - break; - } - } -} diff --git a/src/net/mdns.rs b/src/net/mdns.rs new file mode 100644 index 0000000..7135e7a --- /dev/null +++ b/src/net/mdns.rs @@ -0,0 +1,36 @@ +use alloc::vec; +use embassy_net::udp::UdpSocket; +use embassy_time::{Timer, WithTimeout}; +use sachy_fmt::unwrap; +use sachy_mdns::{GROUP_SOCK_V4, MDNS_PORT, service::MdnsService, state::MdnsAction}; + +pub async fn mdns_loop<'device>(service: &mut MdnsService, udp: &mut UdpSocket<'device>) { + unwrap!(udp.bind(MDNS_PORT)); + + let mut send_buf = vec![0u8; 2048]; + + loop { + if let Some(bytes) = match service.next_action() { + MdnsAction::Announce => service.send_announcement(&mut send_buf), + + MdnsAction::ListenFor { timeout } => udp + .recv_from_with(|buf, _from| service.listen_for_queries(buf, &mut send_buf)) + .with_timeout(timeout) + .await + .ok() + .flatten(), + + MdnsAction::WaitFor { duration } => { + Timer::after(duration).await; + + None + } + } { + if udp.send_to(bytes, GROUP_SOCK_V4).await.is_ok() { + udp.flush().await; + } else { + break; + } + } + } +} diff --git a/src/net/rpc.rs b/src/net/rpc.rs new file mode 100644 index 0000000..0be6296 --- /dev/null +++ b/src/net/rpc.rs @@ -0,0 +1,79 @@ +use alloc::vec; +use embassy_futures::select::select; +use embassy_net::tcp::{TcpReader, TcpSocket, TcpWriter}; +use sachy_fmt::{error, info, unwrap}; + +use crate::{ + constants::HOST_PORT, + rpc::RpcServer, + updates::{NetDataReceiver, UpdateConnection}, +}; + +pub async fn data_loop<'device>(tcp: &mut TcpSocket<'device>, net_data: &NetDataReceiver) { + loop { + UpdateConnection::disconnect(); + + if tcp.accept(HOST_PORT).await.is_err() { + continue; + } + + info!("Connected!"); + UpdateConnection::connect(); + + let (reader, writer) = tcp.split(); + + select(read_loop(reader), write_loop(writer, net_data)).await; + + net_data.clear(); + + tcp.abort(); + tcp.flush().await.ok(); + + info!("DISCONNECT"); + } +} + +async fn read_loop<'device>(mut reader: TcpReader<'device>) { + let mut buf = vec![0u8; 2048]; + + loop { + match reader.read(&mut buf).await { + Ok(0) | Err(_) => break, + Ok(read) => { + if let Ok(req) = striker_proto::receive_request(&mut buf[..read]) + .inspect_err(|e| error!("Proto Error: {}", e)) + { + let Some((resp, resp_tx)) = RpcServer::handle_request(req) + .await + .zip(UpdateConnection::can_update()) + else { + break; + }; + resp_tx.try_send(resp).ok(); + } + } + } + } +} + +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; + } + + if writer.flush().await.is_err() { + break; + } + } +} diff --git a/src/net/sntp.rs b/src/net/sntp.rs new file mode 100644 index 0000000..939d3ed --- /dev/null +++ b/src/net/sntp.rs @@ -0,0 +1,26 @@ +use embassy_net::{dns, udp::UdpSocket}; +use sachy_fmt::{error, unwrap}; +use sachy_sntp::SntpSocket; + +use crate::rtc::GlobalRtc; + +pub async fn sntp_loop<'device>( + udp: &mut UdpSocket<'device>, + stack: embassy_net::Stack<'device>, + rtc: GlobalRtc<'static>, +) { + loop { + let Ok(addr) = stack.dns_query("pool.ntp.org", dns::DnsQueryType::A).await else { + error!("Failed to query DNS for an NTP server. Retrying..."); + continue; + }; + + match udp.resolve_time(addr.as_slice()).await { + Ok(time) => { + unwrap!(rtc.set_rtc_datetime(time).await); + break; + } + Err(e) => error!("Failed to resolve SNTP time: {}", e), + } + } +} diff --git a/src/rpc.rs b/src/rpc.rs index e463cc4..f671097 100644 --- a/src/rpc.rs +++ b/src/rpc.rs @@ -8,23 +8,19 @@ 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 { + 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 a947917..c987205 100644 --- a/src/updates.rs +++ b/src/updates.rs @@ -1,5 +1,5 @@ use embassy_sync::channel::{Channel, Receiver, Sender}; -use striker_proto::{StrikerResponse}; +use striker_proto::StrikerResponse; use crate::{ locks::NetDataLock, diff --git a/src/wifi.rs b/src/wifi.rs index b86d816..4cddbac 100644 --- a/src/wifi.rs +++ b/src/wifi.rs @@ -47,7 +47,7 @@ pub async fn main_loop( control.init(clm).await; control - .set_power_management(cyw43::PowerManagementMode::PowerSave) + .set_power_management(cyw43::PowerManagementMode::ThroughputThrottling) .await; unwrap!(