Skip to content
File

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

rust170 lines
1//! Autonomous SFU session orchestration through the authenticated Worker API.
2use crate::{
3 error::{Error, Result},
4 platform::{self, Http},
5 radio::{Event, Request},
6};
7use radio_core::signaling::{self, Channels, Started};
8use serde::{Serialize, de::DeserializeOwned};
9use std::{
10 ffi::CStr,
11 sync::mpsc::{Receiver, SyncSender, sync_channel},
12 thread,
13 time::Duration,
14};
15 
16#[derive(Clone, Copy)]
17enum Endpoint {
18 Start,
19 Channels,
20 Ready,
21 Heartbeat,
22}
23impl Endpoint {
24 fn path(self) -> &'static CStr {
25 match self {
26 Self::Start => c"start",
27 Self::Channels => c"channels",
28 Self::Ready => c"ready",
29 Self::Heartbeat => c"heartbeat",
30 }
31 }
32}
33 
34struct Client {
35 http: Http,
36}
37impl Client {
38 fn new() -> Result<Self> {
39 Ok(Self {
40 http: Http::open()?,
41 })
42 }
43 fn post(&mut self, endpoint: Endpoint, body: &[u8]) -> Result<(i32, &[u8])> {
44 let (status, response) = self.http.post(endpoint.path(), body)?;
45 platform::log(&format!(
46 "HTTPS {}: status={status}",
47 endpoint.path().to_str().expect("constant ASCII path")
48 ));
49 Ok((status, response))
50 }
51 fn retry<T: DeserializeOwned>(
52 &mut self,
53 endpoint: Endpoint,
54 body: &impl Serialize,
55 ) -> Result<T> {
56 // Reuse the request body across retries, including the initial bootId and offer.
57 let body =
58 serde_json::to_vec(body).map_err(|_| Error::new("request serialization failed"))?;
59 for attempt in 0..3 {
60 let (status, response) = self.post(endpoint, &body)?;
61 if (200..300).contains(&status) {
62 return serde_json::from_slice(response)
63 .map_err(|_| Error::new("invalid Worker response"));
64 }
65 if !signaling::retryable(status) {
66 return Err(Error::native("Worker rejected setup", status));
67 }
68 if attempt < 2 {
69 thread::sleep(Duration::from_secs(1 << attempt));
70 }
71 }
72 Err(Error::new("Worker unavailable after setup retries"))
73 }
74}
75 
76pub(crate) fn run(
77 requests: SyncSender<Request>,
78 events: Receiver<Event>,
79 station: crate::station::Station,
80) -> Result<()> {
81 let mut client = Client::new()?;
82 let (offer, now_playing) = match events.recv_timeout(Duration::from_secs(20)) {
83 Ok(Event::Offer { sdp, now_playing }) => (sdp, now_playing),
84 _ => return Err(Error::new("local offer unavailable")),
85 };
86 let boot_id = format!(
87 "{:08x}{:08x}{:08x}{:08x}",
88 platform::random(),
89 platform::random(),
90 platform::random(),
91 platform::random()
92 );
93 let started: Started = client.retry(
94 Endpoint::Start,
95 &serde_json::json!({
96 "bootId": boot_id, "sessionDescription": {"type": "offer", "sdp": offer},
97 "track": now_playing.track, "nowPlaying": now_playing
98 }),
99 )?;
100 let (identity, answer) = started.validate()?;
101 let (reply, received) = sync_channel(1);
102 requests
103 .try_send(Request::Answer(answer, reply))
104 .map_err(|_| Error::new("radio task unavailable"))?;
105 received
106 .recv_timeout(Duration::from_secs(15))
107 .map_err(|_| Error::new("answer application timed out"))??;
108 match events.recv_timeout(Duration::from_secs(30)) {
109 Ok(Event::Connected) => {}
110 _ => return Err(Error::new("WebRTC connection timed out")),
111 }
112 let channels: Channels = client.retry(Endpoint::Channels, &identity)?;
113 let ids = channels.validate()?;
114 let (reply, received) = sync_channel(1);
115 requests
116 .try_send(Request::Start(ids, reply))
117 .map_err(|_| Error::new("radio task unavailable"))?;
118 received
119 .recv_timeout(Duration::from_secs(15))
120 .map_err(|_| Error::new("playback start timed out"))??;
121 let _: serde_json::Map<String, serde_json::Value> = client.retry(Endpoint::Ready, &identity)?;
122 platform::recovery_attempt(true);
123 platform::log(
124 "ON AIR: autonomous Wi-Fi signaling, audio and data; firmware=rust; USB is optional",
125 );
126 let mut last_success = platform::now_us();
127 let mut next_heartbeat = last_success;
128 let mut next_change_post = last_success;
129 let mut published_revision = None;
130 loop {
131 if events.try_iter().any(|event| matches!(event, Event::Lost)) {
132 return Err(Error::new("WebRTC transport lost"));
133 }
134 thread::sleep(Duration::from_millis(100));
135 let now = platform::now_us();
136 let Some(snapshot) = station.snapshot()? else {
137 continue;
138 };
139 let changed = published_revision != Some(snapshot.revision);
140 if now < next_heartbeat && (!changed || now < next_change_post) {
141 continue;
142 }
143 #[derive(Serialize)]
144 #[serde(rename_all = "camelCase")]
145 struct Heartbeat<'a> {
146 #[serde(flatten)]
147 identity: &'a signaling::Identity,
148 now_playing: &'a radio_core::now_playing::NowPlaying,
149 }
150 let body = serde_json::to_vec(&Heartbeat {
151 identity: &identity,
152 now_playing: &snapshot,
153 })
154 .map_err(|_| Error::new("heartbeat serialization failed"))?;
155 let (status, response) = client.post(Endpoint::Heartbeat, &body)?;
156 let valid = (200..300).contains(&status)
157 && serde_json::from_slice::<serde_json::Map<String, serde_json::Value>>(response)
158 .is_ok();
159 next_heartbeat = platform::now_us() + 5_000_000;
160 next_change_post = platform::now_us() + if valid { 250_000 } else { 1_000_000 };
161 if valid {
162 last_success = platform::now_us();
163 published_revision = Some(snapshot.revision);
164 } else if signaling::needs_recovery(status, platform::now_us().saturating_sub(last_success))
165 {
166 return Err(Error::native("signaling session needs recovery", status));
167 }
168 }
169}