Skip to content
File

Blob: firmware/crates/esp32-radio/src/analysis.rs

rust187 lines
1//! A core-1 task decodes queued copies of music packets. The radio never waits
2//! for analysis: bounded queues and playback epochs discard obsolete work.
3use crate::{
4 error::{Error, Result},
5 platform::{self, Dsp, History},
6};
7use radio_core::{
8 music::{BANDS, MAX_OPUS_BYTES},
9 playback::Due,
10 protocol::SpectrumStatus,
11 spectrum::{Continuity, Tag, Window},
12};
13use std::sync::mpsc::{Receiver, SyncSender, TrySendError, sync_channel};
14 
15struct Job {
16 tag: Tag,
17 length: usize,
18 opus: [u8; MAX_OPUS_BYTES],
19}
20pub(crate) struct Output {
21 pub tag: Tag,
22 pub bands: [u8; BANDS],
23 stats: SpectrumStatus,
24}
25pub(crate) struct Client {
26 input: SyncSender<Job>,
27 output: Receiver<Output>,
28 stats: SpectrumStatus,
29 input_dropped: u64,
30 stale: u64,
31}
32pub(crate) struct Worker {
33 input: Receiver<Job>,
34 output: SyncSender<Output>,
35}
36 
37pub(crate) fn channel() -> (Client, Worker) {
38 let (input, worker_input) = sync_channel(4);
39 let (worker_output, output) = sync_channel(2);
40 (
41 Client {
42 input,
43 output,
44 stats: SpectrumStatus::default(),
45 input_dropped: 0,
46 stale: 0,
47 },
48 Worker {
49 input: worker_input,
50 output: worker_output,
51 },
52 )
53}
54 
55impl Client {
56 pub(crate) fn submit(&mut self, due: Due, opus: &[u8], now_us: u64) -> Result<()> {
57 if opus.is_empty() || opus.len() > MAX_OPUS_BYTES {
58 return Err(Error::new("invalid analysis packet size"));
59 }
60 let mut job = Job {
61 tag: Tag {
62 due,
63 queued_us: now_us,
64 },
65 length: opus.len(),
66 opus: [0; MAX_OPUS_BYTES],
67 };
68 job.opus[..opus.len()].copy_from_slice(opus);
69 match self.input.try_send(job) {
70 Ok(()) => Ok(()),
71 Err(TrySendError::Full(_)) => {
72 self.input_dropped = self.input_dropped.saturating_add(1);
73 Ok(())
74 }
75 Err(TrySendError::Disconnected(_)) => Err(Error::new("analysis task stopped")),
76 }
77 }
78 pub(crate) fn latest(&mut self, epoch: u32, now_us: u64) -> Option<Output> {
79 let mut latest = None;
80 for _ in 0..2 {
81 let Ok(output) = self.output.try_recv() else {
82 break;
83 };
84 self.stats = output.stats;
85 if output.tag.is_fresh(epoch, now_us) {
86 latest = Some(output);
87 } else {
88 self.stale = self.stale.saturating_add(1);
89 }
90 }
91 latest
92 }
93 pub(crate) fn status(&self) -> SpectrumStatus {
94 SpectrumStatus {
95 input_dropped: self.input_dropped,
96 stale: self.stale.saturating_add(self.stats.stale),
97 ..self.stats
98 }
99 }
100}
101 
102impl Worker {
103 pub(crate) fn run(self) -> Result<()> {
104 let mut dsp = Dsp::open()?;
105 let mut history = History::new()?;
106 let mut window =
107 Window::new(history.as_mut()).map_err(|_| Error::new("invalid FFT history"))?;
108 let mut continuity = Continuity::default();
109 let mut stats = SpectrumStatus::default();
110 let mut total_us = 0u64;
111 let mut consecutive_errors = 0;
112 let mut next_stack_sample = 0;
113 platform::log("Live spectrum ready: Rust analysis on core 1, 2048-point ESP-DSP FFT");
114 while let Ok(job) = self.input.recv() {
115 let now = platform::now_us();
116 if !job.tag.is_fresh(job.tag.due.epoch, now) {
117 stats.stale = stats.stale.saturating_add(1);
118 continuity.clear();
119 continue;
120 }
121 if continuity.begin(job.tag.due) {
122 dsp.reset()?;
123 window.clear();
124 stats.resets = stats.resets.saturating_add(1);
125 }
126 let started = platform::now_us();
127 match dsp.decode(&job.opus[..job.length]) {
128 Ok(pcm) => {
129 window
130 .push_stereo(pcm)
131 .map_err(|_| Error::new("invalid PCM shape"))?;
132 consecutive_errors = 0;
133 }
134 Err(_) => {
135 stats.errors = stats.errors.saturating_add(1);
136 consecutive_errors += 1;
137 continuity.clear();
138 window.clear();
139 if consecutive_errors >= 10 {
140 return Err(Error::new("repeated Opus analysis failure"));
141 }
142 continue;
143 }
144 }
145 stats.frames = stats.frames.saturating_add(1);
146 let bands = if job.tag.due.spectrum
147 && window
148 .write_complex(dsp.buffer_mut())
149 .map_err(|_| Error::new("invalid FFT input"))?
150 {
151 dsp.transform()?;
152 Some(
153 window
154 .bands(dsp.buffer())
155 .map_err(|_| Error::new("invalid FFT output"))?,
156 )
157 } else {
158 None
159 };
160 let elapsed = platform::now_us().saturating_sub(started);
161 total_us = total_us.saturating_add(elapsed);
162 stats.mean_us = (total_us / stats.frames).min(u32::MAX as u64) as u32;
163 stats.max_us = stats.max_us.max(elapsed.min(u32::MAX as u64) as u32);
164 if now >= next_stack_sample {
165 stats.stack_free = platform::stack_free();
166 next_stack_sample = now + 5_000_000;
167 }
168 if let Some(bands) = bands {
169 match self.output.try_send(Output {
170 tag: job.tag,
171 bands,
172 stats,
173 }) {
174 Ok(()) => {}
175 Err(TrySendError::Full(_)) => {
176 stats.output_dropped = stats.output_dropped.saturating_add(1)
177 }
178 Err(TrySendError::Disconnected(_)) => {
179 return Err(Error::new("radio task stopped"));
180 }
181 }
182 }
183 }
184 Err(Error::new("analysis input closed"))
185 }
186}