Skip to content
File

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

rust279 lines
1//! One task owns the peer, playback clock, LED, counters, and packet buffers.
2use crate::{
3 error::{Error, Result},
4 platform::{self, Board, MusicStorage},
5};
6use radio_core::{
7 music::{Catalog, MAX_OPUS_BYTES},
8 now_playing::{Message as NowPlayingMessage, NowPlaying, Snapshot},
9 playback::Playback,
10 protocol::{self, Ack, Color, Command, Telemetry},
11 signaling::ChannelIds,
12};
13use radio_webrtc::{Peer, PeerState, SendOutcome, Stream};
14use std::{
15 sync::{
16 Arc,
17 mpsc::{Receiver, SyncSender},
18 },
19 thread,
20 time::Duration,
21};
22 
23pub(crate) enum Request {
24 Answer(String, SyncSender<Result<()>>),
25 Start(ChannelIds, SyncSender<Result<()>>),
26}
27pub(crate) enum Event {
28 Offer {
29 sdp: String,
30 now_playing: NowPlaying,
31 },
32 Connected,
33 Lost,
34}
35 
36// SDP-bearing messages intentionally have no Debug implementation.
37 
38#[derive(Debug, Default)]
39struct Counters {
40 data: u64,
41 rejected: u64,
42 sequence: u64,
43}
44 
45pub(crate) fn run(
46 mut board: Board,
47 requests: Receiver<Request>,
48 events: SyncSender<Event>,
49 mut analysis: crate::analysis::Client,
50 station: crate::station::Station,
51 mut metrics: crate::metrics::Client,
52) -> Result<()> {
53 let mut storage = MusicStorage::open()?;
54 let catalog = Catalog::parse(&mut storage)?;
55 platform::log(&format!(
56 "Music validated: tracks={} bytes={} cache_bytes={} index_bytes={}",
57 catalog.songs().len(),
58 catalog.used_bytes(),
59 storage.cache_bytes(),
60 catalog.index_bytes()
61 ));
62 storage.finish_validation();
63 let mut published = Snapshot::new(&catalog);
64 station.publish(&published.now_playing)?;
65 let mut emitted = false;
66 let mut opus_buffer = [0; MAX_OPUS_BYTES];
67 let certificate = platform::certificate()?;
68 platform::log("DTLS certificate generated on device; starting str0m");
69 let initializing = platform::now_us();
70 let crypto = Arc::new(platform::crypto_provider());
71 let (mut peer, offer) = Peer::open(platform::ipv4()?, certificate, crypto)?;
72 platform::log(&format!(
73 "str0m initialized in {} ms",
74 (platform::now_us() - initializing) / 1000
75 ));
76 events
77 .try_send(Event::Offer {
78 sdp: offer,
79 now_playing: published.now_playing.clone(),
80 })
81 .map_err(|_| Error::new("signaling event queue unavailable"))?;
82 let mut playback = None;
83 let mut led = Color([0, 12, 4]);
84 let mut counters = Counters::default();
85 let mut next_telemetry = 0;
86 let mut next_stack_sample = 0;
87 let mut stack_free = 0;
88 let mut previous = PeerState::Connecting;
89 let mut command_buffer = [0; protocol::COMMAND_LIMIT];
90 let mut json_buffer = Vec::with_capacity(2048);
91 #[cfg(feature = "crypto-profile")]
92 let mut crypto_profile = platform::CryptoProfile::default();
93 loop {
94 if let Ok(request) = requests.try_recv() {
95 match request {
96 Request::Answer(sdp, reply) => {
97 let _ = reply.send(peer.answer(&sdp).map_err(Error::from));
98 }
99 Request::Start(ids, reply) => {
100 let result = if playback.is_some() {
101 Err(Error::new("playback already started"))
102 } else {
103 peer.create_channels(ids).map_err(Error::from)
104 };
105 if result.is_ok() {
106 playback = Playback::playlist(
107 catalog.songs().iter().map(|song| song.frames).collect(),
108 platform::now_us(),
109 );
110 }
111 let _ = reply.send(result);
112 }
113 }
114 }
115 peer.poll()?;
116 let state = peer.state();
117 if state != previous {
118 if state == PeerState::Connected {
119 events
120 .try_send(Event::Connected)
121 .map_err(|_| Error::new("signaling event queue unavailable"))?;
122 platform::log("str0m: WebRTC audio and SCTP connected");
123 } else if state == PeerState::Lost {
124 let _ = events.try_send(Event::Lost);
125 return Err(Error::new("WebRTC transport disconnected"));
126 }
127 previous = state;
128 }
129 if let Some(player) = playback.as_mut().filter(|_| state == PeerState::Connected) {
130 // Bound per-tick work even if a controller floods the reliable stream.
131 let length = peer.command(&mut command_buffer)?;
132 if length > 0 {
133 match Command::parse(&command_buffer[..length]) {
134 Ok(command) => {
135 let mut result = 0;
136 if let Some(color) = command.led {
137 if board.set_led(color).is_ok() {
138 led = color;
139 } else {
140 result = -1;
141 }
142 }
143 if let Some(action) = command.action {
144 player.apply(action);
145 }
146 let ack = Ack {
147 event: "ack",
148 command_id: command.command_id.as_deref(),
149 result,
150 led,
151 paused: player.paused(),
152 };
153 send_json(&mut peer, &ack, &mut json_buffer, &mut counters)?;
154 }
155 Err(_) => counters.rejected = counters.rejected.saturating_add(1),
156 }
157 }
158 let now = platform::now_us();
159 if let Some(due) = player.due(now) {
160 emitted = true;
161 if published.record(&catalog, due)? {
162 station.publish(&published.now_playing)?;
163 send_json(
164 &mut peer,
165 &NowPlayingMessage {
166 event: "nowPlaying",
167 now_playing: &published.now_playing,
168 },
169 &mut json_buffer,
170 &mut counters,
171 )?;
172 }
173 let length = catalog.read_frame(
174 &mut storage,
175 due.track_index as usize,
176 due.index,
177 &mut opus_buffer,
178 )?;
179 let opus = if due.paused {
180 protocol::SILENCE
181 } else {
182 &opus_buffer[..length]
183 };
184 peer.audio(due.pts_ms, opus)?;
185 if !due.paused {
186 analysis.submit(due, opus, now)?;
187 }
188 if due.paused && due.spectrum {
189 send_data(
190 &mut peer,
191 Stream::Spectrum,
192 &protocol::spectrum(due, &[0; 32]),
193 &mut counters,
194 )?;
195 }
196 }
197 if let Some(output) = analysis.latest(player.epoch(), platform::now_us()) {
198 send_data(
199 &mut peer,
200 Stream::Spectrum,
201 &protocol::spectrum(output.tag.due, &output.bands),
202 &mut counters,
203 )?;
204 }
205 if now >= next_telemetry && emitted {
206 next_telemetry = now + 500_000;
207 counters.sequence = counters.sequence.wrapping_add(1);
208 send_json(
209 &mut peer,
210 &NowPlayingMessage {
211 event: "nowPlaying",
212 now_playing: &published.now_playing,
213 },
214 &mut json_buffer,
215 &mut counters,
216 )?;
217 let reading = metrics.latest(now);
218 if now >= next_stack_sample {
219 stack_free = platform::stack_free();
220 next_stack_sample = now + 5_000_000;
221 }
222 let telemetry = Telemetry {
223 event: "telemetry",
224 firmware: "rust",
225 sequence: counters.sequence,
226 uptime_ms: now / 1000,
227 position_ms: published.position_ms,
228 duration_ms: published.duration_ms,
229 playback_revision: published.now_playing.revision,
230 music: protocol::MusicStatus {
231 cache_bytes: storage.cache_bytes() as u32,
232 index_bytes: catalog.index_bytes() as u32,
233 reads: storage.reads,
234 max_read_us: storage.max_read_us,
235 },
236 paused: player.paused(),
237 rssi: reading.and_then(|value| value.rssi),
238 hardware: reading.map(|value| value.hardware),
239 heap: reading
240 .map(|value| value.hardware.internal.free)
241 .unwrap_or_else(platform::heap_free),
242 // Audio send failures end this publisher session.
243 audio_errors: 0,
244 data_errors: counters.data,
245 skipped_frames: player.skipped_frames(),
246 rejected_commands: counters.rejected,
247 dropped_commands: peer.dropped(),
248 stack_free,
249 spectrum: analysis.status(),
250 led,
251 };
252 send_json(&mut peer, &telemetry, &mut json_buffer, &mut counters)?;
253 }
254 #[cfg(feature = "crypto-profile")]
255 crypto_profile.tick(platform::now_us())?;
256 }
257 thread::sleep(Duration::from_millis(1));
258 }
259}
260 
261fn send_json(
262 peer: &mut Peer,
263 value: &impl serde::Serialize,
264 buffer: &mut Vec<u8>,
265 counters: &mut Counters,
266) -> Result<()> {
267 buffer.clear();
268 serde_json::to_writer(&mut *buffer, value)
269 .map_err(|_| Error::new("packet serialization failed"))?;
270 send_data(peer, Stream::Robot, buffer, counters)
271}
272 
273fn send_data(peer: &mut Peer, stream: Stream, bytes: &[u8], counters: &mut Counters) -> Result<()> {
274 if peer.data(stream, bytes)? == SendOutcome::Backpressured {
275 counters.data = counters.data.saturating_add(1);
276 }
277 Ok(())
278}