Skip to content
File

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

rust147 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 response: Vec<u8>,
37}
38impl Client {
39 fn new() -> Result<Self> {
40 Ok(Self {
41 http: Http::open()?,
42 response: vec![0; platform::RESPONSE_LIMIT],
43 })
44 }
45 fn post(&mut self, endpoint: Endpoint, body: &[u8]) -> Result<(i32, usize)> {
46 let (status, length) = self.http.post(endpoint.path(), body, &mut self.response)?;
47 platform::log(&format!(
48 "HTTPS {}: status={status}",
49 endpoint.path().to_str().expect("constant ASCII path")
50 ));
51 Ok((status, length))
52 }
53 fn retry<T: DeserializeOwned>(
54 &mut self,
55 endpoint: Endpoint,
56 body: &impl Serialize,
57 ) -> Result<T> {
58 // Serialize once: retries retain exactly the same bootId and offer.
59 let body =
60 serde_json::to_vec(body).map_err(|_| Error::new("request serialization failed"))?;
61 for attempt in 0..3 {
62 let (status, length) = self.post(endpoint, &body)?;
63 if (200..300).contains(&status) {
64 return serde_json::from_slice(&self.response[..length])
65 .map_err(|_| Error::new("invalid Worker response"));
66 }
67 if !signaling::retryable(status) {
68 return Err(Error::native("Worker rejected setup", status));
69 }
70 if attempt < 2 {
71 thread::sleep(Duration::from_secs(1 << attempt));
72 }
73 }
74 Err(Error::new("Worker unavailable after setup retries"))
75 }
76}
77 
78pub(crate) fn run(requests: SyncSender<Request>, events: Receiver<Event>) -> Result<()> {
79 let mut client = Client::new()?;
80 let offer = match events.recv_timeout(Duration::from_secs(20)) {
81 Ok(Event::Offer(offer)) => offer,
82 _ => return Err(Error::new("local offer unavailable")),
83 };
84 let boot_id = format!(
85 "{:08x}{:08x}{:08x}{:08x}",
86 platform::random(),
87 platform::random(),
88 platform::random(),
89 platform::random()
90 );
91 let started: Started = client.retry(
92 Endpoint::Start,
93 &serde_json::json!({
94 "bootId": boot_id, "sessionDescription": {"type": "offer", "sdp": offer}
95 }),
96 )?;
97 let (identity, answer) = started.validate()?;
98 let (reply, received) = sync_channel(1);
99 requests
100 .try_send(Request::Answer(answer, reply))
101 .map_err(|_| Error::new("radio task unavailable"))?;
102 received
103 .recv_timeout(Duration::from_secs(15))
104 .map_err(|_| Error::new("answer application timed out"))??;
105 match events.recv_timeout(Duration::from_secs(30)) {
106 Ok(Event::Connected) => {}
107 _ => return Err(Error::new("WebRTC connection timed out")),
108 }
109 let channels: Channels = client.retry(Endpoint::Channels, &identity)?;
110 channels.validate()?;
111 let (reply, received) = sync_channel(1);
112 requests
113 .try_send(Request::Start(reply))
114 .map_err(|_| Error::new("radio task unavailable"))?;
115 received
116 .recv_timeout(Duration::from_secs(15))
117 .map_err(|_| Error::new("playback start timed out"))??;
118 let _: serde_json::Map<String, serde_json::Value> = client.retry(Endpoint::Ready, &identity)?;
119 platform::recovery_attempt(true);
120 platform::log(
121 "ON AIR: autonomous Wi-Fi signaling, audio and data; firmware=rust; USB is optional",
122 );
123 let identity =
124 serde_json::to_vec(&identity).map_err(|_| Error::new("identity serialization failed"))?;
125 let mut last_success = platform::now_us();
126 loop {
127 for _ in 0..50 {
128 if events.try_iter().any(|event| matches!(event, Event::Lost)) {
129 return Err(Error::new("WebRTC transport lost"));
130 }
131 thread::sleep(Duration::from_millis(100));
132 }
133 let (status, length) = client.post(Endpoint::Heartbeat, &identity)?;
134 let valid = (200..300).contains(&status)
135 && serde_json::from_slice::<serde_json::Map<String, serde_json::Value>>(
136 &client.response[..length],
137 )
138 .is_ok();
139 if valid {
140 last_success = platform::now_us();
141 } else if signaling::needs_recovery(status, platform::now_us().saturating_sub(last_success))
142 {
143 return Err(Error::native("signaling session needs recovery", status));
144 }
145 }
146}