| 1 | //! Meeting peers through a relay (`crates/relay`): one WebSocket to the room, carrying |
| 2 | //! streams to peers in it, each opened with SPAKE2 and sealed end to end as on a LAN, so the |
| 3 | //! relay sees only who talks to whom, when, and how much. In a code's room every two peers |
| 4 | //! have a stream. In a notebook's room only a host and each guest do; everything else goes |
| 5 | //! to the room's group (`group`), sent once and copied by the relay. A stream or group frame |
| 6 | //! that breaks makes this end join the room again, which ends every stream it had there; |
| 7 | //! peers then meet afresh. |
| 8 | |
| 9 | use super::{ |
| 10 | Event, Hello, OPENING, PATIENCE, Paced, Peer, Pipe, Presence, Relayed as Answer, Room, Shared, |
| 11 | Side, code_parts, group, |
| 12 | transport::{self, Address, Failure, parse}, |
| 13 | wire::kind, |
| 14 | }; |
| 15 | use ::relay::{BROADCAST, GROUP, Notice, SLOT, Verdict, ws}; |
| 16 | use std::{ |
| 17 | collections::{HashMap, HashSet, hash_map::Entry}, |
| 18 | io::{self, BufReader, Read, Write}, |
| 19 | sync::{Arc, Mutex, atomic::Ordering, mpsc}, |
| 20 | thread, |
| 21 | time::{Duration, Instant}, |
| 22 | }; |
| 23 | |
| 24 | /// How often a quiet connection pings the relay. |
| 25 | const KEEPALIVE: Duration = Duration::from_secs(30); |
| 26 | /// A connection that lasted this long was no failure, so the next waits only a second. |
| 27 | const STEADY: Duration = Duration::from_secs(60); |
| 28 | /// The most sent in one message, well under any relay's cap. |
| 29 | const CHUNK: usize = 64 << 10; |
| 30 | /// The largest message read from the relay. |
| 31 | const MOST: usize = 1 << 20; |
| 32 | |
| 33 | type Reader = ws::Reader<BufReader<Box<dyn Read + Send>>>; |
| 34 | |
| 35 | /// Joins `shared`'s room at the relay at `url`, and keeps joining whenever it falls out; |
| 36 | /// `port` is where it listens, to advertise once the relay numbers its code. |
| 37 | pub(super) fn join(shared: &Arc<Shared>, url: &str, port: u16) -> io::Result<()> { |
| 38 | let address = parse(url)?; |
| 39 | let shared = Arc::clone(shared); |
| 40 | thread::Builder::new() |
| 41 | .name("live relay".into()) |
| 42 | .spawn(move || keep(&shared, &address, port))?; |
| 43 | Ok(()) |
| 44 | } |
| 45 | |
| 46 | /// Records how the relay answered, telling `shared`'s events where it changed. |
| 47 | fn answered(shared: &Shared, answer: Answer) { |
| 48 | let mut state = shared.state.lock().unwrap(); |
| 49 | if state.relayed.as_ref() != Some(&answer) { |
| 50 | state.relayed = Some(answer); |
| 51 | drop(state); |
| 52 | (shared.events)(Event::Changed); |
| 53 | } |
| 54 | } |
| 55 | |
| 56 | /// Stays in the room at `address`, joining again whenever the connection ends, until |
| 57 | /// `shared` stops. |
| 58 | fn keep(shared: &Arc<Shared>, address: &Address, port: u16) { |
| 59 | let mut wait = Duration::from_secs(1); |
| 60 | let owner = shared.room.owner(); |
| 61 | // The number of this end's code, asked for again on joining again. |
| 62 | let mut nameplate = owner |
| 63 | .then(|| code_parts(shared.state.lock().unwrap().code.as_deref()?).0) |
| 64 | .flatten(); |
| 65 | while !shared.stopped.load(Ordering::Acquire) { |
| 66 | let began = Instant::now(); |
| 67 | let path = match (owner, shared.tag()) { |
| 68 | (false, Some(tag)) => format!("{}/v1/room/{tag}", address.path), |
| 69 | (false, None) => return, |
| 70 | (true, _) => match nameplate { |
| 71 | Some(number) => format!("{}/v1/claim?nameplate={number}", address.path), |
| 72 | None => format!("{}/v1/claim", address.path), |
| 73 | }, |
| 74 | }; |
| 75 | match connect(shared, address, &path, owner) { |
| 76 | Ok((socket, reader)) => { |
| 77 | shared.state.lock().unwrap().relay = Some(Arc::clone(&socket)); |
| 78 | answered(shared, Answer::Joined); |
| 79 | if !shared.stopped.load(Ordering::Acquire) { |
| 80 | session(shared, &socket, reader, &mut nameplate, port); |
| 81 | } |
| 82 | shared.state.lock().unwrap().relay = None; |
| 83 | socket.hang_up(); |
| 84 | } |
| 85 | Err(Failure::Refused(status, retry)) => { |
| 86 | eprintln!("Live: the relay refused to let this end in ({status})"); |
| 87 | answered(shared, Answer::Refused(status, retry)); |
| 88 | if status == 410 { |
| 89 | return; |
| 90 | } |
| 91 | wait = wait.max(retry.unwrap_or_default()); |
| 92 | } |
| 93 | Err(Failure::Trouble(trouble, error)) => { |
| 94 | eprintln!("Live: no relay at {}: {error}", address.authority); |
| 95 | answered(shared, Answer::Unreachable(trouble)); |
| 96 | } |
| 97 | } |
| 98 | if began.elapsed() >= STEADY { |
| 99 | wait = Duration::from_secs(1); |
| 100 | } |
| 101 | thread::sleep(wait); |
| 102 | wait = (wait * 2).min(PATIENCE); |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | /// Opens the relay's `path` at `address`: the connection, and what reads its messages. |
| 107 | fn connect( |
| 108 | shared: &Shared, |
| 109 | address: &Address, |
| 110 | path: &str, |
| 111 | owner: bool, |
| 112 | ) -> Result<(Arc<Socket>, Reader), Failure> { |
| 113 | let connection = transport::connect(address, path)?; |
| 114 | let reader = ws::Reader::new(BufReader::new(connection.reader), MOST, false); |
| 115 | let group = matches!(shared.room, Room::Notebook(_)).then(|| Group { |
| 116 | keys: group::Keys::new(&shared.secret), |
| 117 | members: Mutex::default(), |
| 118 | out: Mutex::default(), |
| 119 | }); |
| 120 | let socket = Arc::new(Socket { |
| 121 | send: Mutex::new(connection.writer), |
| 122 | close: connection.close, |
| 123 | owner, |
| 124 | links: Mutex::default(), |
| 125 | group, |
| 126 | }); |
| 127 | Ok((socket, reader)) |
| 128 | } |
| 129 | |
| 130 | /// One connection to a relay's room. |
| 131 | pub(super) struct Socket { |
| 132 | send: Mutex<Box<dyn Write + Send>>, |
| 133 | close: Box<dyn Fn() + Send + Sync>, |
| 134 | /// Whether this end claimed the room for its code, and so tells the relay who knew it. |
| 135 | owner: bool, |
| 136 | links: Mutex<Links>, |
| 137 | /// A notebook's room's group. |
| 138 | pub(super) group: Option<Group>, |
| 139 | } |
| 140 | |
| 141 | /// A notebook's room's group: its key, the peers heard in it, and the way to this end's |
| 142 | /// thread sending to it. |
| 143 | pub(super) struct Group { |
| 144 | keys: group::Keys, |
| 145 | /// The peers heard in the group, by slot. |
| 146 | pub(super) members: Mutex<HashMap<u32, group::Member>>, |
| 147 | /// To the thread sending this end's group frames, while connected. |
| 148 | out: Mutex<Option<mpsc::Sender<Out>>>, |
| 149 | } |
| 150 | |
| 151 | /// What this end sends the group next. |
| 152 | pub(super) enum Out { |
| 153 | /// The newest presence, to everyone. |
| 154 | Presence, |
| 155 | /// A message kind and its encoded body, to everyone or to the peers in these slots. |
| 156 | Frame(u16, Vec<u8>, Option<Vec<u32>>), |
| 157 | } |
| 158 | |
| 159 | impl Group { |
| 160 | pub(super) fn send(&self, out: Out) { |
| 161 | if let Some(sender) = &*self.out.lock().unwrap() { |
| 162 | let _ = sender.send(out); |
| 163 | } |
| 164 | } |
| 165 | |
| 166 | /// The slot of the member that is peer `id`. |
| 167 | pub(super) fn slot(&self, id: &[u8; 16]) -> Option<u32> { |
| 168 | let members = self.members.lock().unwrap(); |
| 169 | members |
| 170 | .iter() |
| 171 | .find_map(|(slot, member)| (member.hello()?.peer == *id).then_some(*slot)) |
| 172 | } |
| 173 | } |
| 174 | |
| 175 | /// The streams a socket carries, by the slot of the peer at the other end. |
| 176 | #[derive(Default)] |
| 177 | struct Links { |
| 178 | /// This end's slot. |
| 179 | me: u32, |
| 180 | inboxes: HashMap<u32, mpsc::Sender<Vec<u8>>>, |
| 181 | /// Slots whose stream ended; a peer that comes back has a new slot. |
| 182 | ended: HashSet<u32>, |
| 183 | } |
| 184 | |
| 185 | impl Socket { |
| 186 | fn send(&self, opcode: u8, payload: &[u8]) -> io::Result<()> { |
| 187 | let mut mask = [0; 4]; |
| 188 | getrandom::fill(&mut mask).map_err(|_| io::Error::other("System random source failed"))?; |
| 189 | let frame = ws::frame(opcode, payload, Some(mask)); |
| 190 | self.send.lock().unwrap().write_all(&frame) |
| 191 | } |
| 192 | |
| 193 | pub(super) fn hang_up(&self) { |
| 194 | (self.close)(); |
| 195 | } |
| 196 | |
| 197 | fn forget(&self, slot: u32) { |
| 198 | let mut links = self.links.lock().unwrap(); |
| 199 | links.inboxes.remove(&slot); |
| 200 | links.ended.insert(slot); |
| 201 | } |
| 202 | } |
| 203 | |
| 204 | /// Reads notices and peers' bytes from the relay until the connection ends, opening a stream |
| 205 | /// to each peer in the room. |
| 206 | fn session( |
| 207 | shared: &Arc<Shared>, |
| 208 | socket: &Arc<Socket>, |
| 209 | mut reader: Reader, |
| 210 | nameplate: &mut Option<u32>, |
| 211 | port: u16, |
| 212 | ) { |
| 213 | let mut tag = shared.tag(); |
| 214 | // Pings while the session lasts: dropping `_beat` at its end stops them. |
| 215 | let (_beat, beats) = mpsc::channel::<()>(); |
| 216 | let pinging = Arc::clone(socket); |
| 217 | thread::spawn(move || { |
| 218 | while let Err(mpsc::RecvTimeoutError::Timeout) = beats.recv_timeout(KEEPALIVE) { |
| 219 | if pinging.send(ws::PING, &[]).is_err() { |
| 220 | return; |
| 221 | } |
| 222 | } |
| 223 | }); |
| 224 | if let Some(group) = &socket.group { |
| 225 | let (out, outgoing) = mpsc::channel(); |
| 226 | *group.out.lock().unwrap() = Some(out); |
| 227 | let (shared, socket) = (Arc::clone(shared), Arc::clone(socket)); |
| 228 | thread::spawn(move || speak(&shared, &socket, outgoing)); |
| 229 | } |
| 230 | let notebook = socket.group.is_some(); |
| 231 | while let Ok(message) = reader.read() { |
| 232 | match message { |
| 233 | ws::Message::Text(text) => match text.parse() { |
| 234 | Ok(Notice::Nameplate(number)) => { |
| 235 | *nameplate = Some(number); |
| 236 | tag = Some(format!("code-{number}")); |
| 237 | let super::Room::Code { code, .. } = &shared.room else { |
| 238 | continue; |
| 239 | }; |
| 240 | let Some(code) = super::code::format(number, &code_parts(code).1) else { |
| 241 | continue; |
| 242 | }; |
| 243 | let mut state = shared.state.lock().unwrap(); |
| 244 | if state.code.as_ref() != Some(&code) { |
| 245 | state.code = Some(code); |
| 246 | drop(state); |
| 247 | shared.advertise(port); |
| 248 | (shared.events)(Event::Changed); |
| 249 | } |
| 250 | } |
| 251 | Ok(Notice::Welcome { you, members }) => { |
| 252 | socket.links.lock().unwrap().me = you; |
| 253 | for slot in members.into_iter().filter(|_| !notebook) { |
| 254 | meet(shared, socket, tag.as_deref(), slot, None); |
| 255 | } |
| 256 | } |
| 257 | Ok(Notice::Joined(slot)) if !notebook => { |
| 258 | meet(shared, socket, tag.as_deref(), slot, None); |
| 259 | } |
| 260 | Ok(Notice::Joined(_)) => {} |
| 261 | Ok(Notice::Left(slot)) => { |
| 262 | socket.forget(slot); |
| 263 | let left = (socket.group.as_ref()) |
| 264 | .and_then(|group| group.members.lock().unwrap().remove(&slot)); |
| 265 | if left.is_some() { |
| 266 | (shared.events)(Event::Changed); |
| 267 | } |
| 268 | } |
| 269 | Ok(Notice::Burned) => { |
| 270 | // Coming back, ask for another number. |
| 271 | *nameplate = None; |
| 272 | shared.state.lock().unwrap().burned = true; |
| 273 | (shared.events)(Event::Changed); |
| 274 | } |
| 275 | Err(()) => {} |
| 276 | }, |
| 277 | ws::Message::Binary(data) => match data.split_first_chunk::<SLOT>() { |
| 278 | Some((slot, frame)) if u32::from_be_bytes(*slot) & GROUP != 0 => { |
| 279 | let slot = u32::from_be_bytes(*slot) & !GROUP; |
| 280 | heard(shared, socket, tag.as_deref(), slot, frame); |
| 281 | } |
| 282 | Some((slot, bytes)) => { |
| 283 | let slot = u32::from_be_bytes(*slot); |
| 284 | meet(shared, socket, tag.as_deref(), slot, Some(bytes.to_vec())); |
| 285 | } |
| 286 | None => {} |
| 287 | }, |
| 288 | ws::Message::Ping(payload) => { |
| 289 | let _ = socket.send(ws::PONG, &payload); |
| 290 | } |
| 291 | ws::Message::Pong => {} |
| 292 | ws::Message::Close => break, |
| 293 | } |
| 294 | } |
| 295 | socket.links.lock().unwrap().inboxes.clear(); |
| 296 | if let Some(group) = &socket.group { |
| 297 | *group.out.lock().unwrap() = None; |
| 298 | group.members.lock().unwrap().clear(); |
| 299 | (shared.events)(Event::Changed); |
| 300 | } |
| 301 | } |
| 302 | |
| 303 | /// Sends this end's group frames: its hello, asking everyone for theirs, then its presence |
| 304 | /// at most every `PRESENCE_EVERY`, and frames as they come. |
| 305 | fn speak(shared: &Shared, socket: &Socket, outgoing: mpsc::Receiver<Out>) { |
| 306 | let Some(group) = &socket.group else { |
| 307 | return; |
| 308 | }; |
| 309 | let Ok(mut sealer) = group::Sealer::new(&group.keys) else { |
| 310 | return; |
| 311 | }; |
| 312 | let mut send = |kind: u16, body: &[u8], to: Option<&[u32]>| -> io::Result<()> { |
| 313 | let sealed = sealer.seal(kind, body, to.is_none())?; |
| 314 | let mut message = match to { |
| 315 | None => BROADCAST.to_be_bytes().to_vec(), |
| 316 | Some(slots) => { |
| 317 | let mut message = (GROUP | slots.len() as u32).to_be_bytes().to_vec(); |
| 318 | slots |
| 319 | .iter() |
| 320 | .for_each(|slot| message.extend(slot.to_be_bytes())); |
| 321 | message |
| 322 | } |
| 323 | }; |
| 324 | message.extend(sealed); |
| 325 | socket.send(ws::BINARY, &message) |
| 326 | }; |
| 327 | let hello = minicbor::to_vec(&shared.me).unwrap_or_default(); |
| 328 | if send(kind::HELLO, &hello, None).is_err() { |
| 329 | return; |
| 330 | } |
| 331 | let mut paced = Paced::default(); |
| 332 | paced.changed(); |
| 333 | loop { |
| 334 | let next = match paced.due { |
| 335 | Some(_) => outgoing.recv_timeout(paced.wait(Duration::ZERO)), |
| 336 | None => outgoing |
| 337 | .recv() |
| 338 | .map_err(|_| mpsc::RecvTimeoutError::Disconnected), |
| 339 | }; |
| 340 | let result = match next { |
| 341 | Ok(Out::Presence) => { |
| 342 | paced.changed(); |
| 343 | Ok(()) |
| 344 | } |
| 345 | Ok(Out::Frame(kind, body, to)) => send(kind, &body, to.as_deref()), |
| 346 | Err(mpsc::RecvTimeoutError::Timeout) => Ok(()), |
| 347 | Err(mpsc::RecvTimeoutError::Disconnected) => return, |
| 348 | }; |
| 349 | let result = result.and_then(|()| match paced.due(&shared.state) { |
| 350 | Some(presence) => send( |
| 351 | kind::PRESENCE, |
| 352 | &minicbor::to_vec(&presence).unwrap_or_default(), |
| 353 | None, |
| 354 | ), |
| 355 | None => Ok(()), |
| 356 | }); |
| 357 | if result.is_err() { |
| 358 | return; |
| 359 | } |
| 360 | } |
| 361 | } |
| 362 | |
| 363 | /// Takes in a group frame from the peer in `slot`: its hello, answered with this end's and |
| 364 | /// meeting it where it serves; its presence; or what else it says. |
| 365 | fn heard(shared: &Arc<Shared>, socket: &Arc<Socket>, tag: Option<&str>, slot: u32, frame: &[u8]) { |
| 366 | let Some(group) = &socket.group else { |
| 367 | return; |
| 368 | }; |
| 369 | let opened = { |
| 370 | let mut members = group.members.lock().unwrap(); |
| 371 | let member = match members.entry(slot) { |
| 372 | Entry::Occupied(member) => Ok(member.into_mut()), |
| 373 | Entry::Vacant(vacant) => { |
| 374 | group::Member::new(&group.keys, frame).map(|m| vacant.insert(m)) |
| 375 | } |
| 376 | }; |
| 377 | member.and_then(|member| { |
| 378 | let (kind, body) = member.open(frame)?; |
| 379 | Ok((kind, body, member.hello().cloned())) |
| 380 | }) |
| 381 | }; |
| 382 | let (kind, body, hello) = match opened { |
| 383 | Ok(opened) => opened, |
| 384 | Err(error) => { |
| 385 | eprintln!("Live: the relay broke the room's frames ({error}); meeting again"); |
| 386 | socket.hang_up(); |
| 387 | return; |
| 388 | } |
| 389 | }; |
| 390 | match kind { |
| 391 | kind::HELLO | kind::HELLO_BACK => { |
| 392 | let Ok(hello) = minicbor::decode::<Hello>(&body) else { |
| 393 | return; |
| 394 | }; |
| 395 | if hello.peer == shared.me.peer { |
| 396 | return; |
| 397 | } |
| 398 | let serves = hello.serves.is_some() && shared.me.serves.is_none(); |
| 399 | if let Some(member) = group.members.lock().unwrap().get_mut(&slot) { |
| 400 | member.peer = Some(Peer { |
| 401 | hello: Arc::new(hello), |
| 402 | presence: None, |
| 403 | }); |
| 404 | } |
| 405 | if kind == kind::HELLO { |
| 406 | let presence = minicbor::to_vec(&shared.state.lock().unwrap().presence); |
| 407 | let presence = presence.unwrap_or_default(); |
| 408 | let me = minicbor::to_vec(&shared.me).unwrap_or_default(); |
| 409 | group.send(Out::Frame(kind::HELLO_BACK, me, Some(vec![slot]))); |
| 410 | group.send(Out::Frame(kind::PRESENCE, presence, Some(vec![slot]))); |
| 411 | } |
| 412 | if serves { |
| 413 | meet(shared, socket, tag, slot, None); |
| 414 | } |
| 415 | (shared.events)(Event::Changed); |
| 416 | } |
| 417 | kind::PRESENCE => { |
| 418 | let Ok(presence) = minicbor::decode::<Presence>(&body) else { |
| 419 | return; |
| 420 | }; |
| 421 | let mut members = group.members.lock().unwrap(); |
| 422 | if let Some(peer) = members.get_mut(&slot).and_then(|m| m.peer.as_mut()) { |
| 423 | peer.presence = Some(presence); |
| 424 | drop(members); |
| 425 | (shared.events)(Event::Changed); |
| 426 | } |
| 427 | } |
| 428 | // A peer met directly says the rest there. |
| 429 | kind if kind < 256 => { |
| 430 | if let Some(hello) = hello.filter(|hello| !shared.direct(&hello.peer)) { |
| 431 | (shared.events)(Event::Frame { |
| 432 | from: &hello, |
| 433 | kind, |
| 434 | body: &body, |
| 435 | }); |
| 436 | } |
| 437 | } |
| 438 | _ => {} |
| 439 | } |
| 440 | } |
| 441 | |
| 442 | /// Hands `bytes` from the peer in `slot` to its stream, or opens one. In a code's room this |
| 443 | /// end opens a stream to a peer with a lower slot when told of it, and answers one with a |
| 444 | /// higher slot when its first bytes come; in a notebook's room a guest opens one to its host. |
| 445 | fn meet( |
| 446 | shared: &Arc<Shared>, |
| 447 | socket: &Arc<Socket>, |
| 448 | tag: Option<&str>, |
| 449 | slot: u32, |
| 450 | bytes: Option<Vec<u8>>, |
| 451 | ) { |
| 452 | let mut links = socket.links.lock().unwrap(); |
| 453 | if let Some(inbox) = links.inboxes.get(&slot) { |
| 454 | if let Some(bytes) = bytes { |
| 455 | let _ = inbox.send(bytes); |
| 456 | } |
| 457 | return; |
| 458 | } |
| 459 | let notebook = socket.group.is_some(); |
| 460 | let side = match bytes { |
| 461 | None if notebook || slot < links.me => Side::Initiator, |
| 462 | Some(_) if (notebook && shared.me.serves.is_some()) || (!notebook && slot > links.me) => { |
| 463 | Side::Responder |
| 464 | } |
| 465 | _ => return, |
| 466 | }; |
| 467 | let Some(tag) = tag.filter(|_| !links.ended.contains(&slot)) else { |
| 468 | return; |
| 469 | }; |
| 470 | let (inbox, arriving) = mpsc::channel(); |
| 471 | if let Some(bytes) = bytes { |
| 472 | let _ = inbox.send(bytes); |
| 473 | } |
| 474 | links.inboxes.insert(slot, inbox); |
| 475 | let pipe = Arc::new(Relayed { |
| 476 | socket: Arc::clone(socket), |
| 477 | slot, |
| 478 | inbox: Mutex::new(Inbox { |
| 479 | arriving, |
| 480 | chunk: Vec::new(), |
| 481 | at: 0, |
| 482 | }), |
| 483 | timeout: Mutex::new(OPENING), |
| 484 | }); |
| 485 | let shared = Arc::clone(shared); |
| 486 | let tag = tag.to_owned(); |
| 487 | thread::spawn(move || shared.run(pipe, side, &tag)); |
| 488 | } |
| 489 | |
| 490 | /// The stream to one peer through the relay. |
| 491 | struct Relayed { |
| 492 | socket: Arc<Socket>, |
| 493 | slot: u32, |
| 494 | inbox: Mutex<Inbox>, |
| 495 | timeout: Mutex<Duration>, |
| 496 | } |
| 497 | |
| 498 | struct Inbox { |
| 499 | arriving: mpsc::Receiver<Vec<u8>>, |
| 500 | chunk: Vec<u8>, |
| 501 | at: usize, |
| 502 | } |
| 503 | |
| 504 | impl Pipe for Relayed { |
| 505 | fn read(&self, buffer: &mut [u8]) -> io::Result<usize> { |
| 506 | let mut inbox = self.inbox.lock().unwrap(); |
| 507 | while inbox.at == inbox.chunk.len() { |
| 508 | let timeout = *self.timeout.lock().unwrap(); |
| 509 | inbox.chunk = match inbox.arriving.recv_timeout(timeout) { |
| 510 | Ok(chunk) => chunk, |
| 511 | Err(mpsc::RecvTimeoutError::Timeout) => return Err(io::ErrorKind::TimedOut.into()), |
| 512 | Err(mpsc::RecvTimeoutError::Disconnected) => return Ok(0), |
| 513 | }; |
| 514 | inbox.at = 0; |
| 515 | } |
| 516 | let at = inbox.at; |
| 517 | let length = buffer.len().min(inbox.chunk.len() - at); |
| 518 | buffer[..length].copy_from_slice(&inbox.chunk[at..at + length]); |
| 519 | inbox.at += length; |
| 520 | Ok(length) |
| 521 | } |
| 522 | |
| 523 | fn write(&self, bytes: &[u8]) -> io::Result<usize> { |
| 524 | if self.socket.links.lock().unwrap().ended.contains(&self.slot) { |
| 525 | return Err(io::ErrorKind::BrokenPipe.into()); |
| 526 | } |
| 527 | let length = bytes.len().min(CHUNK); |
| 528 | let message = [&self.slot.to_be_bytes()[..], &bytes[..length]].concat(); |
| 529 | self.socket.send(ws::BINARY, &message)?; |
| 530 | Ok(length) |
| 531 | } |
| 532 | |
| 533 | fn set_read_timeout(&self, timeout: Duration) -> io::Result<()> { |
| 534 | *self.timeout.lock().unwrap() = timeout; |
| 535 | Ok(()) |
| 536 | } |
| 537 | |
| 538 | fn shutdown(&self) { |
| 539 | self.socket.forget(self.slot); |
| 540 | } |
| 541 | |
| 542 | fn direct(&self) -> bool { |
| 543 | false |
| 544 | } |
| 545 | |
| 546 | fn met(&self, met: bool) { |
| 547 | if self.socket.owner { |
| 548 | let verdict = if met { Verdict::Met } else { Verdict::Failed }(self.slot); |
| 549 | let _ = self.socket.send(ws::TEXT, verdict.to_string().as_bytes()); |
| 550 | } |
| 551 | } |
| 552 | |
| 553 | fn broken(&self) { |
| 554 | // A deliberate shutdown can wake the reader with EOF. |
| 555 | if !self.socket.links.lock().unwrap().ended.contains(&self.slot) { |
| 556 | self.socket.hang_up(); |
| 557 | } |
| 558 | } |
| 559 | } |
| 560 | |
| 561 | #[cfg(test)] |
| 562 | mod tests { |
| 563 | use super::*; |
| 564 | |
| 565 | #[test] |
| 566 | fn relay_addresses_parse() { |
| 567 | let address = parse("wss://live.example.net").unwrap(); |
| 568 | assert_eq!( |
| 569 | address, |
| 570 | Address { |
| 571 | tls: true, |
| 572 | authority: "live.example.net".into(), |
| 573 | host: "live.example.net".into(), |
| 574 | port: 443, |
| 575 | path: String::new(), |
| 576 | } |
| 577 | ); |
| 578 | let address = parse("ws://[::1]:7650/snowbound/").unwrap(); |
| 579 | assert_eq!((address.host.as_str(), address.port), ("::1", 7650)); |
| 580 | assert_eq!(address.path, "/snowbound"); |
| 581 | assert!(parse("https://live.example.net").is_err()); |
| 582 | assert!(parse("ws://:80").is_err()); |
| 583 | } |
| 584 | } |