1//! Browser Live Share: event-driven relay connections carrying the same sealed streams.
2
3pub use ::relay::code;
4#[path = "group.rs"]
5mod group;
6#[path = "model.rs"]
7mod model;
8#[path = "share.rs"]
9pub mod share;
10#[path = "wire.rs"]
11pub mod wire;
12use model::hex;
13pub use model::{Event, Peer, Reach, Relayed, Room, Trouble};
14pub use wire::{Caret, Guid, Hello, Presence, Spot};
15
16use minicbor::Encode;
17use std::{
18 collections::BTreeMap,
19 io,
20 sync::{
21 Arc, Mutex, Weak,
22 atomic::{AtomicBool, Ordering},
23 },
24 time::Duration,
25};
26use wasm_bindgen::prelude::*;
27use wire::{Side, kind};
28
29#[wasm_bindgen(module = "/src/live/web.js")]
30extern "C" {
31 #[wasm_bindgen(catch, js_name = liveConnect)]
32 fn connect(url: &str, message: &JsValue, closed: &JsValue) -> Result<u32, JsValue>;
33 #[wasm_bindgen(catch, js_name = liveSend)]
34 fn send(id: u32, bytes: &[u8]) -> Result<(), JsValue>;
35 #[wasm_bindgen(js_name = liveClose)]
36 fn close(id: u32);
37}
38
39pub struct Live(Arc<Shared>);
40
41struct Shared {
42 me: Hello,
43 room: Room,
44 url: String,
45 events: Box<dyn Fn(Event) + Send + Sync>,
46 stopped: AtomicBool,
47 state: Mutex<State>,
48}
49
50struct State {
51 socket: u32,
52 slot: u32,
53 relayed: Relayed,
54 failed: u32,
55 outdated: Option<u16>,
56 burned: bool,
57 streams: BTreeMap<u32, Stream>,
58 members: BTreeMap<u32, group::Member>,
59 sealer: Option<group::Sealer>,
60 presence: Presence,
61 scheduled: bool,
62}
63
64struct Stream {
65 bytes: Vec<u8>,
66 opening: Option<wire::Opening>,
67 receive: Option<wire::Sealer>,
68 line: Line,
69 peer: Option<Peer>,
70}
71
72#[derive(Clone)]
73pub struct Line {
74 shared: Weak<Shared>,
75 slot: u32,
76 sealer: Arc<Mutex<Option<wire::Sealer>>>,
77}
78
79impl Line {
80 pub fn send(&self, kind: u16, body: &impl Encode<()>) -> io::Result<()> {
81 let shared = self.shared.upgrade().ok_or(io::ErrorKind::NotConnected)?;
82 let mut bytes = Vec::new();
83 self.sealer
84 .lock()
85 .unwrap()
86 .as_mut()
87 .ok_or(io::ErrorKind::NotConnected)?
88 .send(&mut bytes, kind, body)?;
89 shared.stream(self.slot, &bytes)
90 }
91
92 pub fn hang_up(&self, reason: &str) {
93 let _ = self.send(
94 kind::BYE,
95 &wire::Bye {
96 reason: reason.into(),
97 },
98 );
99 }
100}
101
102#[derive(Clone)]
103pub struct Sender(Weak<Shared>);
104
105impl Sender {
106 pub fn send(&self, kind: u16, body: &impl Encode<()>, to: Option<&[[u8; 16]]>) {
107 if let Some(shared) = self.0.upgrade() {
108 match to {
109 None => {
110 let _ = shared.group(kind, body, None);
111 }
112 Some(to) => {
113 let slots: Vec<_> = shared
114 .state
115 .lock()
116 .unwrap()
117 .members
118 .iter()
119 .filter(|(_, m)| {
120 m.peer.as_ref().is_some_and(|p| to.contains(&p.hello.peer))
121 })
122 .map(|(slot, _)| *slot)
123 .collect();
124 for slot in slots {
125 let _ = shared.group(kind, body, Some(slot));
126 }
127 }
128 }
129 }
130 }
131}
132
133impl Live {
134 pub fn start(
135 me: Hello,
136 room: &Room,
137 _: Option<Reach>,
138 relay: Option<&str>,
139 events: impl Fn(Event) + Send + Sync + 'static,
140 ) -> io::Result<Self> {
141 let tag = room.tag().ok_or(io::ErrorKind::InvalidInput)?;
142 let relay = relay
143 .ok_or(io::ErrorKind::NotConnected)?
144 .trim_end_matches('/');
145 if !relay.starts_with("wss://") && !relay.starts_with("ws://") {
146 return Err(io::ErrorKind::InvalidInput.into());
147 }
148 let shared = Arc::new(Shared {
149 me,
150 room: room.clone(),
151 url: format!("{relay}/v1/room/{tag}"),
152 events: Box::new(events),
153 stopped: AtomicBool::new(false),
154 state: Mutex::new(State {
155 socket: 0,
156 slot: 0,
157 relayed: Relayed::Unknown,
158 failed: 0,
159 outdated: None,
160 burned: false,
161 streams: BTreeMap::new(),
162 members: BTreeMap::new(),
163 sealer: None,
164 presence: Presence::default(),
165 scheduled: false,
166 }),
167 });
168 shared.connect()?;
169 let weak = Arc::downgrade(&shared);
170 crate::task::spawn("live ping", move || async move {
171 let (_, wait) = crate::task::channel();
172 loop {
173 crate::task::wait(&wait, Some(Duration::from_secs(15))).await;
174 let Some(shared) = weak
175 .upgrade()
176 .filter(|s| !s.stopped.load(Ordering::Acquire))
177 else {
178 break;
179 };
180 let lines: Vec<_> = shared
181 .state
182 .lock()
183 .unwrap()
184 .streams
185 .values()
186 .map(|s| s.line.clone())
187 .collect();
188 for line in lines {
189 let _ = line.send(kind::PING, &());
190 }
191 if matches!(shared.room, Room::Notebook(_)) {
192 let _ = shared.group(kind::PING, &(), None);
193 }
194 }
195 })?;
196 Ok(Self(shared))
197 }
198
199 pub fn code(&self) -> Option<String> {
200 match &self.0.room {
201 Room::Code { code, .. } => Some(code.clone()),
202 _ => None,
203 }
204 }
205 pub fn relayed(&self) -> Relayed {
206 self.0.state.lock().unwrap().relayed.clone()
207 }
208 pub fn failed(&self) -> u32 {
209 self.0.state.lock().unwrap().failed
210 }
211 pub fn other_version(&self) -> Option<u16> {
212 self.0.state.lock().unwrap().outdated
213 }
214 pub fn burned(&self) -> bool {
215 self.0.state.lock().unwrap().burned
216 }
217 pub fn sender(&self) -> Sender {
218 Sender(Arc::downgrade(&self.0))
219 }
220 pub fn peers(&self) -> Vec<Peer> {
221 let state = self.0.state.lock().unwrap();
222 let mut peers = BTreeMap::new();
223 for peer in state
224 .members
225 .values()
226 .filter_map(|m| m.peer.as_ref())
227 .chain(state.streams.values().filter_map(|s| s.peer.as_ref()))
228 {
229 peers.insert(peer.hello.peer, peer.clone());
230 }
231 peers.into_values().collect()
232 }
233 pub fn line(&self, peer: &[u8; 16]) -> Option<Line> {
234 self.0
235 .state
236 .lock()
237 .unwrap()
238 .streams
239 .values()
240 .find(|s| s.peer.as_ref().is_some_and(|p| &p.hello.peer == peer))
241 .map(|s| s.line.clone())
242 }
243 pub fn set_presence(&self, presence: Presence) {
244 let mut state = self.0.state.lock().unwrap();
245 if state.presence == presence {
246 return;
247 }
248 state.presence = presence;
249 if std::mem::replace(&mut state.scheduled, true) {
250 return;
251 }
252 drop(state);
253 let weak = Arc::downgrade(&self.0);
254 let _ = crate::task::spawn("live presence", move || async move {
255 let (_, wait) = crate::task::channel();
256 crate::task::wait(&wait, Some(Duration::from_millis(100))).await;
257 if let Some(shared) = weak.upgrade() {
258 let presence = {
259 let mut state = shared.state.lock().unwrap();
260 state.scheduled = false;
261 state.presence.clone()
262 };
263 let _ = shared.group(kind::PRESENCE, &presence, None);
264 }
265 });
266 }
267 pub fn leave(self, reason: &str) {
268 let lines: Vec<_> = self
269 .0
270 .state
271 .lock()
272 .unwrap()
273 .streams
274 .values()
275 .map(|s| s.line.clone())
276 .collect();
277 for line in lines {
278 line.hang_up(reason);
279 }
280 }
281}
282
283impl Drop for Live {
284 fn drop(&mut self) {
285 self.0.stopped.store(true, Ordering::Release);
286 close(self.0.state.lock().unwrap().socket);
287 }
288}
289
290impl Shared {
291 fn connect(self: &Arc<Self>) -> io::Result<()> {
292 let weak = Arc::downgrade(self);
293 let message = Closure::<dyn FnMut(JsValue)>::new(move |data: JsValue| {
294 if let Some(shared) = weak.upgrade()
295 && let Err(error) = shared.heard(data)
296 {
297 if let Some(version) = error
298 .get_ref()
299 .and_then(|e| e.downcast_ref::<wire::Version>())
300 {
301 shared.state.lock().unwrap().outdated = Some(version.0);
302 } else if error.kind() == io::ErrorKind::InvalidData {
303 shared.state.lock().unwrap().failed += 1;
304 }
305 shared.disconnected();
306 }
307 })
308 .into_js_value();
309 let weak = Arc::downgrade(self);
310 let closed = Closure::<dyn FnMut()>::new(move || {
311 if let Some(shared) = weak.upgrade() {
312 shared.disconnected();
313 }
314 })
315 .into_js_value();
316 let socket = connect(&self.url, &message, &closed).map_err(js_error)?;
317 self.state.lock().unwrap().socket = socket;
318 Ok(())
319 }
320
321 fn disconnected(self: &Arc<Self>) {
322 let peers = {
323 let mut state = self.state.lock().unwrap();
324 close(state.socket);
325 let peers: Vec<_> = state
326 .streams
327 .values()
328 .filter_map(|s| s.peer.as_ref().map(|p| p.hello.clone()))
329 .collect();
330 state.streams.clear();
331 state.members.clear();
332 state.sealer = None;
333 state.relayed = Relayed::Unreachable(Trouble::Other);
334 peers
335 };
336 for hello in peers {
337 (self.events)(Event::Left(&hello));
338 }
339 (self.events)(Event::Changed);
340 let weak = Arc::downgrade(self);
341 let _ = crate::task::spawn("live reconnect", move || async move {
342 let (_, wait) = crate::task::channel();
343 crate::task::wait(&wait, Some(Duration::from_secs(2))).await;
344 if let Some(shared) = weak
345 .upgrade()
346 .filter(|s| !s.stopped.load(Ordering::Acquire))
347 {
348 let _ = shared.connect();
349 }
350 });
351 }
352
353 fn stream(&self, slot: u32, bytes: &[u8]) -> io::Result<()> {
354 let socket = self.state.lock().unwrap().socket;
355 for chunk in bytes.chunks(64 << 10) {
356 send(socket, &[&slot.to_be_bytes()[..], chunk].concat()).map_err(js_error)?;
357 }
358 Ok(())
359 }
360
361 fn group(&self, kind: u16, body: &impl Encode<()>, to: Option<u32>) -> io::Result<()> {
362 let body = minicbor::to_vec(body).map_err(io::Error::other)?;
363 let (socket, sealed) = {
364 let mut state = self.state.lock().unwrap();
365 let sealed = state
366 .sealer
367 .as_mut()
368 .ok_or(io::ErrorKind::NotConnected)?
369 .seal(kind, &body, to.is_none())?;
370 (state.socket, sealed)
371 };
372 let mut bytes = match to {
373 Some(slot) => [&(::relay::GROUP | 1).to_be_bytes()[..], &slot.to_be_bytes()].concat(),
374 None => ::relay::BROADCAST.to_be_bytes().to_vec(),
375 };
376 bytes.extend(sealed);
377 send(socket, &bytes).map_err(js_error)
378 }
379
380 fn meet(self: &Arc<Self>, slot: u32) -> io::Result<()> {
381 if self.state.lock().unwrap().streams.contains_key(&slot) {
382 return Ok(());
383 }
384 let tag = self.room.tag().ok_or(io::ErrorKind::InvalidInput)?;
385 let (opening, bytes) = wire::Opening::new(Side::Initiator, &tag, &self.room.secret())?;
386 let line = Line {
387 shared: Arc::downgrade(self),
388 slot,
389 sealer: Arc::new(Mutex::new(None)),
390 };
391 self.state.lock().unwrap().streams.insert(
392 slot,
393 Stream {
394 bytes: Vec::new(),
395 opening: Some(opening),
396 receive: None,
397 line,
398 peer: None,
399 },
400 );
401 self.stream(
402 slot,
403 &[&(bytes.len() as u32).to_be_bytes()[..], &bytes].concat(),
404 )
405 }
406
407 fn heard(self: &Arc<Self>, data: JsValue) -> io::Result<()> {
408 if let Some(text) = data.as_string() {
409 let notice = text
410 .parse::<::relay::Notice>()
411 .map_err(|_| io::ErrorKind::InvalidData)?;
412 match notice {
413 ::relay::Notice::Welcome { you, members } => {
414 {
415 let mut state = self.state.lock().unwrap();
416 state.slot = you;
417 state.relayed = Relayed::Joined;
418 if matches!(self.room, Room::Notebook(_)) {
419 state.sealer =
420 Some(group::Sealer::new(&group::Keys::new(&self.room.secret()))?);
421 }
422 }
423 if matches!(self.room, Room::Code { .. }) {
424 for slot in members.into_iter().filter(|slot| *slot < you) {
425 self.meet(slot)?;
426 }
427 } else {
428 self.group(kind::HELLO, &self.me, None)?;
429 let presence = self.state.lock().unwrap().presence.clone();
430 self.group(kind::PRESENCE, &presence, None)?;
431 }
432 }
433 ::relay::Notice::Left(slot) => {
434 let hello = {
435 let mut state = self.state.lock().unwrap();
436 state.members.remove(&slot);
437 state
438 .streams
439 .remove(&slot)
440 .and_then(|s| s.peer.map(|p| p.hello))
441 };
442 if let Some(hello) = hello {
443 (self.events)(Event::Left(&hello));
444 }
445 }
446 ::relay::Notice::Burned => self.state.lock().unwrap().burned = true,
447 _ => {}
448 }
449 (self.events)(Event::Changed);
450 return Ok(());
451 }
452 let bytes = js_sys::Uint8Array::new(&data).to_vec();
453 let (slot, bytes) = bytes
454 .split_first_chunk::<4>()
455 .ok_or(io::ErrorKind::InvalidData)?;
456 let slot = u32::from_be_bytes(*slot);
457 if slot & ::relay::GROUP != 0 {
458 return self.heard_group(slot & !::relay::GROUP, bytes);
459 }
460 self.heard_stream(slot, bytes)
461 }
462
463 fn heard_group(self: &Arc<Self>, slot: u32, frame: &[u8]) -> io::Result<()> {
464 let (kind, body, hello) = {
465 let mut state = self.state.lock().unwrap();
466 let member = match state.members.entry(slot) {
467 std::collections::btree_map::Entry::Occupied(e) => e.into_mut(),
468 std::collections::btree_map::Entry::Vacant(e) => e.insert(group::Member::new(
469 &group::Keys::new(&self.room.secret()),
470 frame,
471 )?),
472 };
473 let (kind, body) = member.open(frame)?;
474 (kind, body, member.hello().cloned())
475 };
476 match kind {
477 kind::HELLO | kind::HELLO_BACK => {
478 let hello: Hello =
479 minicbor::decode(&body).map_err(|_| io::ErrorKind::InvalidData)?;
480 if hello.peer == self.me.peer {
481 return Ok(());
482 }
483 let serves = hello.serves.is_some() && self.me.serves.is_none();
484 self.state
485 .lock()
486 .unwrap()
487 .members
488 .get_mut(&slot)
489 .unwrap()
490 .peer = Some(Peer {
491 hello: Arc::new(hello),
492 presence: None,
493 });
494 if kind == kind::HELLO {
495 self.group(kind::HELLO_BACK, &self.me, Some(slot))?;
496 let presence = self.state.lock().unwrap().presence.clone();
497 self.group(kind::PRESENCE, &presence, Some(slot))?;
498 }
499 if serves {
500 self.meet(slot)?;
501 }
502 (self.events)(Event::Changed);
503 }
504 kind::PRESENCE => {
505 let presence = minicbor::decode(&body).map_err(|_| io::ErrorKind::InvalidData)?;
506 if let Some(peer) = self
507 .state
508 .lock()
509 .unwrap()
510 .members
511 .get_mut(&slot)
512 .and_then(|m| m.peer.as_mut())
513 {
514 peer.presence = Some(presence);
515 }
516 (self.events)(Event::Changed);
517 }
518 kind if kind < 256 => {
519 if let Some(hello) = hello {
520 (self.events)(Event::Frame {
521 from: &hello,
522 kind,
523 body: &body,
524 });
525 }
526 }
527 _ => {}
528 }
529 Ok(())
530 }
531
532 fn heard_stream(self: &Arc<Self>, slot: u32, bytes: &[u8]) -> io::Result<()> {
533 {
534 let mut state = self.state.lock().unwrap();
535 let stream = state
536 .streams
537 .get_mut(&slot)
538 .ok_or(io::ErrorKind::InvalidData)?;
539 if stream.bytes.len() + bytes.len() > (16 << 20) + 4 {
540 return Err(io::ErrorKind::InvalidData.into());
541 }
542 stream.bytes.extend_from_slice(bytes);
543 }
544 loop {
545 let event = {
546 let mut state = self.state.lock().unwrap();
547 let stream = state.streams.get_mut(&slot).unwrap();
548 let Some(length) = stream.bytes.get(..4) else {
549 break;
550 };
551 let length = u32::from_be_bytes(length.try_into().unwrap()) as usize;
552 if length > 16 << 20 {
553 return Err(io::ErrorKind::InvalidData.into());
554 }
555 if stream.bytes.len() < length + 4 {
556 break;
557 }
558 let bytes: Vec<_> = stream.bytes.drain(..length + 4).collect();
559 if let Some(opening) = stream.opening.take() {
560 let (send, receive) = opening.finish(&bytes[4..])?;
561 *stream.line.sealer.lock().unwrap() = Some(send);
562 stream.receive = Some(receive);
563 Some((stream.line.clone(), None, kind::HELLO, Vec::new()))
564 } else {
565 let (kind, body) = stream
566 .receive
567 .as_mut()
568 .ok_or(io::ErrorKind::InvalidData)?
569 .receive(&mut io::Cursor::new(bytes))?;
570 if stream.peer.is_none() {
571 if kind != kind::HELLO {
572 return Err(io::ErrorKind::InvalidData.into());
573 }
574 let hello = minicbor::decode::<Hello>(&body)
575 .map_err(|_| io::ErrorKind::InvalidData)?;
576 stream.peer = Some(Peer {
577 hello: Arc::new(hello),
578 presence: None,
579 });
580 }
581 Some((
582 stream.line.clone(),
583 stream.peer.as_ref().map(|p| p.hello.clone()),
584 kind,
585 body,
586 ))
587 }
588 };
589 if let Some((line, hello, kind, body)) = event {
590 match hello {
591 None => line.send(kind::HELLO, &self.me)?,
592 Some(hello) if kind == kind::HELLO => (self.events)(Event::Met(&hello, &line)),
593 Some(hello) => (self.events)(Event::Frame {
594 from: &hello,
595 kind,
596 body: &body,
597 }),
598 }
599 }
600 }
601 Ok(())
602 }
603}
604
605fn js_error(error: JsValue) -> io::Error {
606 io::Error::new(
607 io::ErrorKind::NotConnected,
608 error
609 .as_string()
610 .unwrap_or_else(|| "Relay disconnected".into()),
611 )
612}