File
Blob: firmware/crates/esp32-radio/src/radio.rs
| 1 | //! One task owns the peer, playback clock, LED, counters, and packet buffers. |
| 2 | use crate::{ |
| 3 | error::{Error, Result}, |
| 4 | platform::{self, Board, MusicStorage}, |
| 5 | }; |
| 6 | use 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 | }; |
| 13 | use radio_webrtc::{Peer, PeerState, SendOutcome, Stream}; |
| 14 | use std::{ |
| 15 | sync::{ |
| 16 | Arc, |
| 17 | mpsc::{Receiver, SyncSender}, |
| 18 | }, |
| 19 | thread, |
| 20 | time::Duration, |
| 21 | }; |
| 22 | |
| 23 | pub(crate) enum Request { |
| 24 | Answer(String, SyncSender<Result<()>>), |
| 25 | Start(ChannelIds, SyncSender<Result<()>>), |
| 26 | } |
| 27 | pub(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)] |
| 39 | struct Counters { |
| 40 | data: u64, |
| 41 | rejected: u64, |
| 42 | sequence: u64, |
| 43 | } |
| 44 | |
| 45 | pub(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 | |
| 261 | fn 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 | |
| 273 | fn 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 | } |