From bf2d93ad5a531553b0e2c3c7153701c6715534ff Mon Sep 17 00:00:00 2001 From: Sachymetsu Date: Tue, 23 Dec 2025 17:03:10 +0000 Subject: [PATCH] Timestamped detection loop --- src/adc.rs | 4 +-- src/detector.rs | 62 +++++++++++++++++++++++++++------------- src/detector/analysis.rs | 6 +++- src/detector/data.rs | 13 +++++---- src/main.rs | 2 +- src/net.rs | 6 +++- src/pwm.rs | 3 +- src/rtc.rs | 19 ++++++++++-- src/state.rs | 2 +- src/updates.rs | 6 +++- src/utils.rs | 6 ++-- 11 files changed, 90 insertions(+), 39 deletions(-) diff --git a/src/adc.rs b/src/adc.rs index c71c230..ab039ac 100644 --- a/src/adc.rs +++ b/src/adc.rs @@ -1,8 +1,8 @@ use embassy_rp::{ + Peri, adc::{self, Adc, AdcPin, Async, Config}, dma, peripherals::ADC, - Peri, }; use embassy_sync::{blocking_mutex::raw::RawMutex, mutex::Mutex}; @@ -50,7 +50,7 @@ impl<'device, M: RawMutex, T: dma::Channel> AdcDriver<'device, M, T> { /// it tries again until the sampling succeeds. pub async fn sample(&self, buf: &mut [u16]) { // Gain a lock to the inner ADC state, then do a read. - self.inner.lock().await.read_many(buf.as_mut()).await; + self.inner.lock().await.read_many(buf).await; } pub async fn sample_average(&self, buf: &mut [u16]) -> u16 { diff --git a/src/detector.rs b/src/detector.rs index 2d0db8d..3c3f963 100644 --- a/src/detector.rs +++ b/src/detector.rs @@ -1,10 +1,11 @@ mod analysis; mod data; -use core::cell::{Cell, RefCell, RefMut}; +use core::cell::Cell; use alloc::vec::Vec; -use defmt::{info, unwrap}; +use chrono::TimeDelta; +use sachy_fmt::{info, unwrap}; use embassy_futures::select::select3; use embassy_rp::peripherals::DMA_CH1; @@ -17,25 +18,32 @@ use embassy_time::{Duration, Instant, Ticker, Timer}; use crate::{ adc::AdcDriver, pwm::PwmDriver, + rtc::GlobalRtc, state::{DEVICE_STATE, DeviceState}, updates::{NET_CHANNEL, NetDataSender}, - utils::{static_alloc, try_buffer, try_static_block_vecs}, + utils::{static_alloc, try_buffer, try_static_timestamped_block_vecs}, }; #[embassy_executor::task] pub async fn detector_task( adc: AdcDriver<'static, NoopRawMutex, DMA_CH1>, pwm: PwmDriver<'static>, + rtc: GlobalRtc<'static>, ) { + while !rtc.is_running().await { + info!("Waiting for RTC to start running"); + Timer::after_secs(2).await; + } let net_data = NET_CHANNEL.sender(); info!("Allocating detector resources"); let blocks = unwrap!( - try_static_block_vecs(4, 512), + try_static_timestamped_block_vecs(4, 512), "Couldn't allocate block buffers" ); - let buf_channel: &mut ZChannel> = static_alloc(ZChannel::new(blocks)); + let buf_channel: &mut ZChannel)> = + static_alloc(ZChannel::new(blocks)); let (mut sender, mut receiver) = buf_channel.split(); @@ -52,9 +60,9 @@ pub async fn detector_task( let average = detector.tune(samples.as_mut_slice()).await; select3( - detector.sample(&mut sender), + detector.sample(&mut sender, rtc), detector.analyse(&mut receiver, average, &mut peaks), - detector.tick(), + detector.tick(rtc), ) .await; } @@ -64,7 +72,7 @@ struct Detector<'device> { adc: AdcDriver<'device, NoopRawMutex, DMA_CH1>, pwm: PwmDriver<'device>, state: DetectorState, - net_data: RefCell, + net_data: NetDataSender, } #[derive(Debug)] @@ -93,7 +101,7 @@ impl<'device> Detector<'device> { Self { adc, pwm, - net_data: RefCell::new(net_data), + net_data, state: DetectorState::default(), } } @@ -135,9 +143,18 @@ impl Detector<'_> { self.adc.sample_average(samples).await } - async fn sample(&self, data: &mut Sender<'static, NoopRawMutex, Vec>) { + async fn sample( + &self, + data: &mut Sender<'static, NoopRawMutex, (i64, Vec)>, + rtc: GlobalRtc<'static>, + ) { + let time = unwrap!(rtc.get_timestamp().await, "Unable to get timestamp").and_utc(); + let now = Instant::now(); + loop { - let buf = data.send().await; + let (timestamp, buf) = data.send().await; + let elapsed = now.elapsed().as_micros() as i64; + *timestamp = (time + TimeDelta::microseconds(elapsed)).timestamp(); self.adc.sample(buf).await; data.send_done(); } @@ -145,14 +162,14 @@ impl Detector<'_> { async fn analyse( &self, - data: &mut Receiver<'static, NoopRawMutex, Vec>, + data: &mut Receiver<'static, NoopRawMutex, (i64, Vec)>, mut average: u16, peaks: &mut Vec, ) { loop { peaks.clear(); - let buf = data.receive().await; + let (timestamp, buf) = data.receive().await; let new_avg = analysis::analyse_buffer_by_stepped_windows(buf.as_slice(), average, peaks); @@ -165,7 +182,13 @@ impl Detector<'_> { .set(self.state.strikes.get().saturating_add(32)); if let Some(net_data) = self.get_data_channel() { - data::transmit_strike(buf.as_slice(), peaks.as_slice(), average, net_data); + data::transmit_strike( + *timestamp, + buf.as_slice(), + peaks.as_slice(), + average, + net_data, + ); } info!( @@ -184,7 +207,7 @@ impl Detector<'_> { } } - async fn tick(&self) { + async fn tick(&self, rtc: GlobalRtc<'static>) { let mut inactive: Option = None; let mut interval = Ticker::every(Duration::from_secs(1)); @@ -212,7 +235,8 @@ impl Detector<'_> { } None => { if let Some(net_data) = self.get_data_channel() { - data::transmit_level_update(decay, net_data); + let now = unwrap!(rtc.get_timestamp().await, "Unable to get timestamp"); + data::transmit_level_update(now.and_utc().timestamp(), decay, net_data); } } _ => continue, @@ -220,11 +244,9 @@ impl Detector<'_> { } } - fn get_data_channel(&self) -> Option> { + fn get_data_channel(&self) -> Option<&NetDataSender> { let can_send_data = DEVICE_STATE.lock(|x| x.get() == DeviceState::Connected); - can_send_data - .then(|| self.net_data.try_borrow_mut().ok()) - .flatten() + can_send_data.then_some(&self.net_data) } } diff --git a/src/detector/analysis.rs b/src/detector/analysis.rs index 377621a..ab83a47 100644 --- a/src/detector/analysis.rs +++ b/src/detector/analysis.rs @@ -6,7 +6,11 @@ use crate::BLOCK_SIZE; /// needs to be a consistent drop over more than the threshold over a minimum amount of time. const BLIP_THRESHOLD: u16 = 5; -pub(super) fn analyse_buffer_by_stepped_windows(buf: &[u16], average: u16, peaks: &mut Vec) -> u16 { +pub(super) fn analyse_buffer_by_stepped_windows( + buf: &[u16], + average: u16, + peaks: &mut Vec, +) -> u16 { const CHUNK_SIZE: usize = BLOCK_SIZE / 32; const CHUNK_STEP: usize = CHUNK_SIZE / 2; diff --git a/src/detector/data.rs b/src/detector/data.rs index 2459085..77ae8b3 100644 --- a/src/detector/data.rs +++ b/src/detector/data.rs @@ -1,21 +1,24 @@ -use core::cell::RefMut; - use crate::updates::NetDataSender; -pub(super) fn transmit_level_update(warn_level: u16, net_data: RefMut<'_, NetDataSender>) { +pub(super) fn transmit_level_update(timestamp: i64, warn_level: u16, net_data: &NetDataSender) { net_data - .try_send(crate::updates::Update::Warning(warn_level)) + .try_send(crate::updates::Update::Warning { + timestamp, + level: warn_level, + }) .ok(); } pub(super) fn transmit_strike( + timestamp: i64, samples: &[u16], peaks: &[u16], average: u16, - net_data: RefMut<'_, NetDataSender>, + net_data: &NetDataSender, ) { net_data .try_send(crate::updates::Update::Strike { + timestamp, peaks: peaks.to_vec(), samples: samples.to_vec(), average, diff --git a/src/main.rs b/src/main.rs index 0335a92..91d618e 100644 --- a/src/main.rs +++ b/src/main.rs @@ -86,7 +86,7 @@ fn main() -> ! { EXECUTOR1.init_with(Executor::new).run(|spawner| { info!("Spawning Detector task"); - spawner.must_spawn(detector::detector_task(adc, pwm)); + spawner.must_spawn(detector::detector_task(adc, pwm, rtc)); }) }); diff --git a/src/net.rs b/src/net.rs index 4774290..d054291 100644 --- a/src/net.rs +++ b/src/net.rs @@ -10,7 +10,11 @@ use sachy_fmt::{error, unwrap}; use sachy_mdns::{GROUP_ADDR_V4, GROUP_SOCK_V4, MDNS_PORT, MdnsAction, MdnsService, Service}; use sachy_sntp::SntpSocket; -use crate::{rtc::GlobalRtc, state::DEVICE_STATE, updates::{NetDataReceiver, NET_CHANNEL}}; +use crate::{ + rtc::GlobalRtc, + state::DEVICE_STATE, + updates::{NET_CHANNEL, NetDataReceiver}, +}; #[embassy_executor::task] pub async fn udp_stack(stack: embassy_net::Stack<'static>, rtc: GlobalRtc<'static>) { diff --git a/src/pwm.rs b/src/pwm.rs index 90ca4d2..944c9c5 100644 --- a/src/pwm.rs +++ b/src/pwm.rs @@ -1,5 +1,6 @@ use embassy_rp::{ - pwm::{ChannelBPin, Config as PwmConfig, Pwm, Slice}, Peri, + Peri, + pwm::{ChannelBPin, Config as PwmConfig, Pwm, Slice}, }; pub struct PwmDriver<'device> { diff --git a/src/rtc.rs b/src/rtc.rs index c892ec3..cd5d504 100644 --- a/src/rtc.rs +++ b/src/rtc.rs @@ -1,6 +1,6 @@ #![allow(dead_code)] -use chrono::{Datelike, Timelike}; +use chrono::{Datelike, NaiveDate, NaiveDateTime, NaiveTime, Timelike}; use embassy_rp::{ peripherals, rtc::{DateTime, DateTimeError, DayOfWeek, Rtc, RtcError}, @@ -30,10 +30,23 @@ impl GlobalRtc<'_> { rtc.is_running() } - pub async fn get_timestamp(&self) -> Result { + pub async fn get_timestamp(&self) -> Result { let rtc = self.0.lock().await; - rtc.now() + let now = rtc.now()?; + + let date = unwrap!(NaiveDate::from_ymd_opt( + now.year as i32, + now.month as u32, + now.day as u32 + )); + let time = unwrap!(NaiveTime::from_hms_opt( + now.hour as u32, + now.minute as u32, + now.second as u32 + )); + + Ok(NaiveDateTime::new(date, time)) } pub async fn set_rtc_datetime(&self, timestamp: SntpTimestamp) -> Result<(), RtcError> { diff --git a/src/state.rs b/src/state.rs index a183d7f..e4276b4 100644 --- a/src/state.rs +++ b/src/state.rs @@ -11,4 +11,4 @@ pub static DEVICE_STATE: Mutex> = pub enum DeviceState { Disconnected, Connected, -} \ No newline at end of file +} diff --git a/src/updates.rs b/src/updates.rs index 07c72b5..129d4d6 100644 --- a/src/updates.rs +++ b/src/updates.rs @@ -10,8 +10,12 @@ pub type NetDataReceiver = Receiver<'static, NetDataLock, Update, 8>; #[derive(Debug, serde::Serialize)] pub enum Update { - Warning(u16), + Warning { + timestamp: i64, + level: u16, + }, Strike { + timestamp: i64, peaks: Vec, samples: Vec, average: u16, diff --git a/src/utils.rs b/src/utils.rs index ba03610..40ec1ba 100644 --- a/src/utils.rs +++ b/src/utils.rs @@ -14,10 +14,10 @@ pub fn try_buffer(capacity: usize) -> Result, TryReserveError Ok(buffer) } -pub fn try_static_block_vecs( +pub fn try_static_timestamped_block_vecs( block_num: usize, block_capacity: usize, -) -> Result<&'static mut [Vec], TryReserveError> { +) -> Result<&'static mut [(i64, Vec)], TryReserveError> { let mut blocks = Vec::new(); blocks.try_reserve_exact(block_num)?; @@ -25,7 +25,7 @@ pub fn try_static_block_vecs( for _ in 0..block_num { let block = try_buffer(block_capacity)?; - blocks.push(block); + blocks.push((0, block)); } Ok(blocks.leak()) -- 2.51.2