authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-02 22:54:39-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-03 05:26:36-07:00
log8d1b56375d7289dbb12b8ac754cf673adcb8e19e
tree87b4958745edafe43549043d52f86f3aab91633b
parent3dff91e5b8938209f2a7b38e8323227ba4a12a71
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

feat: Live Share reaches its relay through proxies and networks that refuse WebSockets

The relay connection, and only it and update downloads, goes through the proxy HTTPS_PROXY or ALL_PROXY names (short of NO_PROXY), else the system's: macOS's network settings, Windows's Internet settings, GNOME's. It asks an HTTP proxy to CONNECT, with the name and password the proxy URL holds, and trusts the certificates the system trusts, so a proxy that inspects HTTPS works once the system trusts its authority. SMB never goes through it. Where something on the way refuses the WebSocket, the app joins the same room with ?poll=1 instead: GET /v1/poll/<id> waits for what is to come and POST sends, each body a run of WebSocket frames, so the relay hears the same messages either way. This needs the new relay; deploy it. Why the relay can't be reached is told apart and said with what to try: its name not found, nothing answering (a firewall), no answer in time, the proxy not there, the proxy asking for a password, the proxy refusing (its status), a certificate the system doesn't trust, or the network blocking the relay. Tested through a local proxy that asks for a password, one that refuses WebSockets inside its tunnels (the guest then joins and publishes by polling), a proxy that isn't there and a relay name that doesn't resolve. Assisted-by: claude-opus-5.5

14 files changed, 1592 insertions(+), 421 deletions(-)

crates/notebook/Cargo.toml+1-1
...@@ -52,7 +52,7 @@ rsqlite-vfs = "0.1.1"...@@ -52,7 +52,7 @@ rsqlite-vfs = "0.1.1"
52nix = { version = "0.31", default-features = false, features = ["fs"] }52nix = { version = "0.31", default-features = false, features = ["fs"] }
5353
54[target.'cfg(windows)'.dependencies]54[target.'cfg(windows)'.dependencies]
55windows-sys = { version = "0.61", features = ["Win32_Storage_FileSystem"] }55windows-sys = { version = "0.61", features = ["Win32_Storage_FileSystem", "Win32_System_Registry"] }
5656
57[dev-dependencies]57[dev-dependencies]
58libc = "0.2"58libc = "0.2"
crates/notebook/src/live.rs+29-2
...@@ -8,8 +8,11 @@...@@ -8,8 +8,11 @@
8//! from scratch.8//! from scratch.
99
10pub use ::relay::code;10pub use ::relay::code;
11pub mod proxy;
11mod relay;12mod relay;
12pub mod share;13pub mod share;
14mod transport;
15pub use transport::Trouble;
13pub mod wire;16pub mod wire;
14pub use wire::{Caret, Guid, Hello, Presence, Spot};17pub use wire::{Caret, Guid, Hello, Presence, Spot};
1518
...@@ -156,8 +159,8 @@ pub enum Relayed {...@@ -156,8 +159,8 @@ pub enum Relayed {
156 Joined,159 Joined,
157 /// Answered with this HTTP status, and how long it asked to wait.160 /// Answered with this HTTP status, and how long it asked to wait.
158 Refused(u16, Option<Duration>),161 Refused(u16, Option<Duration>),
159 /// Not reached.162 /// Not reached, and why.
160 Unreachable,163 Unreachable(Trouble),
161}164}
162165
163/// Presence on the network while it lives; dropping it leaves.166/// Presence on the network while it lives; dropping it leaves.
...@@ -755,6 +758,30 @@ fn discover(...@@ -755,6 +758,30 @@ fn discover(
755 Ok(daemon)758 Ok(daemon)
756}759}
757760
761/// `text` with `%XX` escapes decoded, as a URL's name and password.
762fn decode(text: &str) -> String {
763 let bytes = text.as_bytes();
764 let mut decoded = Vec::with_capacity(bytes.len());
765 let mut at = 0;
766 while at < bytes.len() {
767 match text
768 .get(at + 1..at + 3)
769 .filter(|_| bytes[at] == b'%')
770 .and_then(|hex| u8::from_str_radix(hex, 16).ok())
771 {
772 Some(byte) => {
773 decoded.push(byte);
774 at += 3;
775 }
776 None => {
777 decoded.push(bytes[at]);
778 at += 1;
779 }
780 }
781 }
782 String::from_utf8_lossy(&decoded).into_owned()
783}
784
758fn hex(bytes: &[u8]) -> String {785fn hex(bytes: &[u8]) -> String {
759 bytes.iter().map(|byte| format!("{byte:02x}")).collect()786 bytes.iter().map(|byte| format!("{byte:02x}")).collect()
760}787}
crates/notebook/src/live/proxy.rs created+290
...@@ -0,0 +1,290 @@
1//! The proxy a relay connection, and only a relay connection, goes through: `HTTPS_PROXY`
2//! (`HTTP_PROXY` for `ws://`), then `ALL_PROXY`, each in lower case too, short of `NO_PROXY`;
3//! else the system's: macOS's network settings, Windows's Internet settings, GNOME's. Only
4//! HTTP proxies, reached by `CONNECT`, with a name and password where the URL has them; a
5//! proxy auto-configuration script is not read.
6
7use base64::Engine;
8use std::sync::Mutex;
9
10#[derive(Clone, Debug, PartialEq, Eq)]
11pub struct Proxy {
12 pub host: String,
13 pub port: u16,
14 /// The name and password `Proxy-Authorization` sends.
15 pub credentials: Option<(String, String)>,
16}
17
18impl Proxy {
19 /// `http://[name:password@]host[:port][/]`, or `host:port`; none for another scheme.
20 pub fn parse(url: &str) -> Option<Self> {
21 let url = url.trim();
22 let rest = match url.split_once("://") {
23 Some(("http" | "https", rest)) => rest,
24 Some(_) => return None,
25 None => url,
26 };
27 let rest = rest.split('/').next()?;
28 let (credentials, authority) = match rest.rsplit_once('@') {
29 Some((credentials, authority)) => {
30 let (name, password) = credentials.split_once(':').unwrap_or((credentials, ""));
31 (
32 Some((crate::live::decode(name), crate::live::decode(password))),
33 authority,
34 )
35 }
36 None => (None, rest),
37 };
38 let (host, port) = match authority.rsplit_once(':') {
39 Some((host, port)) if !port.contains(']') => (host, port.parse().ok()?),
40 _ => (authority, 80),
41 };
42 let host = host.trim_start_matches('[').trim_end_matches(']');
43 (!host.is_empty()).then(|| Self {
44 host: host.to_owned(),
45 port,
46 credentials,
47 })
48 }
49
50 /// The `Proxy-Authorization` header's value, where it has credentials.
51 pub fn authorization(&self) -> Option<String> {
52 let (name, password) = self.credentials.as_ref()?;
53 let token = base64::engine::general_purpose::STANDARD.encode(format!("{name}:{password}"));
54 Some(format!("Basic {token}"))
55 }
56}
57
58impl std::fmt::Display for Proxy {
59 fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
60 write!(f, "{}:{}", self.host, self.port)
61 }
62}
63
64/// A proxy `use_proxy` set in place of what the environment and the system say.
65static CHOSEN: Mutex<Option<Option<Proxy>>> = Mutex::new(None);
66
67/// Uses `proxy`, or no proxy with `Some(None)`, for every relay connection from now on; `None`
68/// goes back to the environment's and the system's.
69pub fn use_proxy(proxy: Option<Option<Proxy>>) {
70 *CHOSEN
71 .lock()
72 .unwrap_or_else(|poisoned| poisoned.into_inner()) = proxy;
73}
74
75/// The proxy a connection to `host` goes through, over TLS where `tls`.
76pub fn for_host(host: &str, tls: bool) -> Option<Proxy> {
77 if let Some(chosen) = CHOSEN.lock().unwrap_or_else(|p| p.into_inner()).clone() {
78 return chosen;
79 }
80 let variable = |name: &str| {
81 [name.to_owned(), name.to_lowercase()]
82 .into_iter()
83 .find_map(|name| std::env::var(name).ok().filter(|value| !value.is_empty()))
84 };
85 let named = [if tls { "HTTPS_PROXY" } else { "HTTP_PROXY" }, "ALL_PROXY"]
86 .into_iter()
87 .find_map(variable);
88 match named {
89 Some(url) => {
90 let bypass = variable("NO_PROXY").unwrap_or_default();
91 (!bypassed(host, bypass.split(','))).then(|| Proxy::parse(&url))?
92 }
93 None => system(host, tls),
94 }
95}
96
97/// Whether `host` is one of `list`'s: a name, a suffix after a dot, or `*` for every host.
98fn bypassed<'a>(host: &str, list: impl Iterator<Item = &'a str>) -> bool {
99 let host = host.to_ascii_lowercase();
100 list.map(|entry| {
101 entry
102 .trim()
103 .trim_start_matches("*.")
104 .trim_start_matches('.')
105 })
106 .filter(|entry| !entry.is_empty())
107 .any(|entry| {
108 let entry = entry.to_ascii_lowercase();
109 entry == "*" || host == entry || host.ends_with(&format!(".{entry}"))
110 })
111}
112
113/// What `scutil --proxy` says of the network settings in use.
114#[cfg(target_os = "macos")]
115fn system(host: &str, tls: bool) -> Option<Proxy> {
116 let output = std::process::Command::new("/usr/sbin/scutil")
117 .arg("--proxy")
118 .output()
119 .ok()?;
120 let text = String::from_utf8(output.stdout).ok()?;
121 let value = |key: &str| {
122 text.lines().find_map(|line| {
123 let (name, value) = line.split_once(" : ")?;
124 (name.trim() == key).then(|| value.trim().to_owned())
125 })
126 };
127 let kind = if tls { "HTTPS" } else { "HTTP" };
128 if value(&format!("{kind}Enable")).as_deref() != Some("1") {
129 return None;
130 }
131 let exceptions: Vec<String> = text
132 .lines()
133 .skip_while(|line| !line.contains("ExceptionsList"))
134 .skip(1)
135 .take_while(|line| !line.contains('}'))
136 .filter_map(|line| Some(line.split_once(" : ")?.1.trim().to_owned()))
137 .collect();
138 if bypassed(host, exceptions.iter().map(String::as_str)) {
139 return None;
140 }
141 Some(Proxy {
142 host: value(&format!("{kind}Proxy"))?,
143 port: value(&format!("{kind}Port"))?.parse().ok()?,
144 credentials: None,
145 })
146}
147
148/// What Windows's Internet settings say, as WinHTTP and Internet Explorer read them for this
149/// user: `ProxyServer`, as `host:port` or `https=host:port;http=...`, where `ProxyEnable`.
150#[cfg(windows)]
151fn system(host: &str, tls: bool) -> Option<Proxy> {
152 let key = r"Software\Microsoft\Windows\CurrentVersion\Internet Settings";
153 if registry::dword(key, "ProxyEnable")? == 0 {
154 return None;
155 }
156 let server = registry::string(key, "ProxyServer")?;
157 let overrides = registry::string(key, "ProxyOverride").unwrap_or_default();
158 if bypassed(
159 host,
160 overrides.split(';').filter(|entry| *entry != "<local>"),
161 ) {
162 return None;
163 }
164 let scheme = if tls { "https=" } else { "http=" };
165 let server = match server.contains('=') {
166 true => server
167 .split(';')
168 .find_map(|entry| entry.trim().strip_prefix(scheme))?,
169 false => server.as_str(),
170 };
171 Proxy::parse(server)
172}
173
174#[cfg(windows)]
175#[allow(unsafe_code)]
176mod registry {
177 use windows_sys::Win32::System::Registry::{
178 HKEY_CURRENT_USER, RRF_RT_REG_DWORD, RRF_RT_REG_SZ, RegGetValueW,
179 };
180
181 fn wide(text: &str) -> Vec<u16> {
182 text.encode_utf16().chain([0]).collect()
183 }
184
185 pub fn dword(key: &str, name: &str) -> Option<u32> {
186 let (key, name) = (wide(key), wide(name));
187 let mut value = 0u32;
188 let mut size = 4u32;
189 // SAFETY: the key and name are NUL-terminated and outlive the call; the value is four
190 // bytes, as `size` says.
191 let result = unsafe {
192 RegGetValueW(
193 HKEY_CURRENT_USER,
194 key.as_ptr(),
195 name.as_ptr(),
196 RRF_RT_REG_DWORD,
197 std::ptr::null_mut(),
198 (&mut value as *mut u32).cast(),
199 &mut size,
200 )
201 };
202 (result == 0).then_some(value)
203 }
204
205 pub fn string(key: &str, name: &str) -> Option<String> {
206 let (key, name) = (wide(key), wide(name));
207 let mut buffer = vec![0u16; 2048];
208 let mut size = (buffer.len() * 2) as u32;
209 // SAFETY: as in `dword`, with `size` the buffer's length in bytes.
210 let result = unsafe {
211 RegGetValueW(
212 HKEY_CURRENT_USER,
213 key.as_ptr(),
214 name.as_ptr(),
215 RRF_RT_REG_SZ,
216 std::ptr::null_mut(),
217 buffer.as_mut_ptr().cast(),
218 &mut size,
219 )
220 };
221 if result != 0 {
222 return None;
223 }
224 let length = buffer
225 .iter()
226 .position(|unit| *unit == 0)
227 .unwrap_or(buffer.len());
228 Some(String::from_utf16_lossy(&buffer[..length]))
229 }
230}
231
232/// What GNOME's proxy settings say, where they are manual.
233#[cfg(not(any(target_os = "macos", windows)))]
234fn system(host: &str, tls: bool) -> Option<Proxy> {
235 let get = |schema: &str, key: &str| {
236 let output = std::process::Command::new("gsettings")
237 .args(["get", schema, key])
238 .output()
239 .ok()?;
240 let text = String::from_utf8(output.stdout).ok()?;
241 Some(text.trim().trim_matches('\'').to_owned())
242 };
243 if get("org.gnome.system.proxy", "mode")? != "manual" {
244 return None;
245 }
246 let ignored = get("org.gnome.system.proxy", "ignore-hosts").unwrap_or_default();
247 let ignored = ignored.trim_matches(['[', ']']).replace('\'', "");
248 if bypassed(host, ignored.split(',')) {
249 return None;
250 }
251 let schema = if tls {
252 "org.gnome.system.proxy.https"
253 } else {
254 "org.gnome.system.proxy.http"
255 };
256 Some(Proxy {
257 host: get(schema, "host").filter(|host| !host.is_empty())?,
258 port: get(schema, "port")?
259 .parse()
260 .ok()
261 .filter(|port| *port != 0)?,
262 credentials: None,
263 })
264}
265
266#[cfg(test)]
267mod tests {
268 use super::*;
269
270 #[test]
271 fn proxies_read_as_written() {
272 assert_eq!(
273 Proxy::parse("http://ada:p%40ss@proxy.example:3128/"),
274 Some(Proxy {
275 host: "proxy.example".into(),
276 port: 3128,
277 credentials: Some(("ada".into(), "p@ss".into())),
278 })
279 );
280 assert_eq!(Proxy::parse("proxy:8080").unwrap().port, 8080);
281 assert_eq!(Proxy::parse("http://[::1]:8080").unwrap().host, "::1");
282 assert!(Proxy::parse("socks5://proxy:1080").is_none());
283 let proxy = Proxy::parse("http://ada:secret@proxy:1").unwrap();
284 assert_eq!(proxy.authorization().unwrap(), "Basic YWRhOnNlY3JldA==");
285 let list = || ["localhost", ".internal", "*.corp.example"].into_iter();
286 assert!(bypassed("localhost", list()) && bypassed("files.internal", list()));
287 assert!(bypassed("a.corp.example", list()) && !bypassed("relay.example", list()));
288 assert!(bypassed("anything", ["*"].into_iter()));
289 }
290}
crates/notebook/src/live/relay.rs+16-213
...@@ -3,23 +3,21 @@...@@ -3,23 +3,21 @@
3//! the relay sees only who talks to whom, when, and how much. A stream that breaks makes this3//! the relay sees only who talks to whom, when, and how much. A stream that breaks makes this
4//! end join the room again, which ends every stream it had there; peers then meet afresh.4//! end join the room again, which ends every stream it had there; peers then meet afresh.
55
6use super::{Event, OPENING, PATIENCE, Pipe, Relayed as Answer, Shared, Side, code_parts};6use super::{
7 Event, OPENING, PATIENCE, Pipe, Relayed as Answer, Shared, Side, code_parts,
8 transport::{self, Address, Failure, parse},
9};
7use ::relay::{Notice, SLOT, Verdict, ws};10use ::relay::{Notice, SLOT, Verdict, ws};
8use base64::Engine;
9use rustls::{ClientConfig, ClientConnection, RootCertStore, pki_types::ServerName};
10use std::{11use std::{
11 collections::{HashMap, HashSet},12 collections::{HashMap, HashSet},
12 io::{self, BufReader, Read, Write},13 io::{self, BufReader, Read, Write},
13 net::{Shutdown, TcpStream, ToSocketAddrs},14 sync::{Arc, Mutex, atomic::Ordering, mpsc},
14 sync::{Arc, Mutex, OnceLock, atomic::Ordering, mpsc},
15 thread,15 thread,
16 time::{Duration, Instant},16 time::{Duration, Instant},
17};17};
1818
19const CONNECT: Duration = Duration::from_secs(10);19/// How often a quiet connection pings the relay.
20/// How often a quiet connection pings the relay, and how long it waits to hear anything.
21const KEEPALIVE: Duration = Duration::from_secs(30);20const KEEPALIVE: Duration = Duration::from_secs(30);
22const QUIET: Duration = Duration::from_secs(75);
23/// A connection that lasted this long was no failure, so the next waits only a second.21/// A connection that lasted this long was no failure, so the next waits only a second.
24const STEADY: Duration = Duration::from_secs(60);22const STEADY: Duration = Duration::from_secs(60);
25/// The most sent in one message, well under any relay's cap.23/// The most sent in one message, well under any relay's cap.
...@@ -40,54 +38,6 @@ pub(super) fn join(shared: &Arc<Shared>, url: &str, port: u16) -> io::Result<()>...@@ -40,54 +38,6 @@ pub(super) fn join(shared: &Arc<Shared>, url: &str, port: u16) -> io::Result<()>
40 Ok(())38 Ok(())
41}39}
4240
43/// A relay's address: `ws://` or `wss://`, a host, and the path the relay's `/v1/` follows.
44#[derive(Debug, PartialEq)]
45struct Address {
46 tls: bool,
47 /// The host and port as the URL gave them, for the `Host` header.
48 authority: String,
49 host: String,
50 port: u16,
51 path: String,
52}
53
54fn parse(url: &str) -> io::Result<Address> {
55 let bad = || io::Error::new(io::ErrorKind::InvalidInput, format!("Not a relay: {url}"));
56 let (tls, rest) = match url.split_once("://") {
57 Some(("wss", rest)) => (true, rest),
58 Some(("ws", rest)) => (false, rest),
59 _ => return Err(bad()),
60 };
61 let (authority, path) = rest.split_at(rest.find('/').unwrap_or(rest.len()));
62 let (host, port) = match authority.rsplit_once(':') {
63 Some((host, port)) if !port.contains(']') => (host, port.parse().map_err(|_| bad())?),
64 _ => (authority, if tls { 443 } else { 80 }),
65 };
66 let host = host.trim_start_matches('[').trim_end_matches(']');
67 if host.is_empty() {
68 return Err(bad());
69 }
70 Ok(Address {
71 tls,
72 authority: authority.into(),
73 host: host.into(),
74 port,
75 path: path.trim_end_matches('/').into(),
76 })
77}
78
79enum Failure {
80 Network(io::Error),
81 /// The relay answered, but with this HTTP status and how long to wait.
82 Refused(u16, Option<Duration>),
83}
84
85impl From<io::Error> for Failure {
86 fn from(error: io::Error) -> Self {
87 Failure::Network(error)
88 }
89}
90
91/// Records how the relay answered, telling `shared`'s events where it changed.41/// Records how the relay answered, telling `shared`'s events where it changed.
92fn answered(shared: &Shared, answer: Answer) {42fn answered(shared: &Shared, answer: Answer) {
93 let mut state = shared.state.lock().unwrap();43 let mut state = shared.state.lock().unwrap();
...@@ -135,9 +85,9 @@ fn keep(shared: &Arc<Shared>, address: &Address, port: u16) {...@@ -135,9 +85,9 @@ fn keep(shared: &Arc<Shared>, address: &Address, port: u16) {
135 }85 }
136 wait = wait.max(retry.unwrap_or_default());86 wait = wait.max(retry.unwrap_or_default());
137 }87 }
138 Err(Failure::Network(error)) => {88 Err(Failure::Trouble(trouble, error)) => {
139 eprintln!("Live: no relay at {}: {error}", address.authority);89 eprintln!("Live: no relay at {}: {error}", address.authority);
140 answered(shared, Answer::Unreachable);90 answered(shared, Answer::Unreachable(trouble));
141 }91 }
142 }92 }
143 if began.elapsed() >= STEADY {93 if began.elapsed() >= STEADY {
...@@ -148,170 +98,23 @@ fn keep(shared: &Arc<Shared>, address: &Address, port: u16) {...@@ -148,170 +98,23 @@ fn keep(shared: &Arc<Shared>, address: &Address, port: u16) {
148 }98 }
149}99}
150100
151/// Opens a WebSocket to `path` at `address`: the connection, and its reading half.101/// Opens the relay's `path` at `address`: the connection, and what reads its messages.
152fn connect(address: &Address, path: &str, owner: bool) -> Result<(Arc<Socket>, Reader), Failure> {102fn connect(address: &Address, path: &str, owner: bool) -> Result<(Arc<Socket>, Reader), Failure> {
153 let target = (address.host.as_str(), address.port)103 let connection = transport::connect(address, path)?;
154 .to_socket_addrs()?104 let reader = ws::Reader::new(BufReader::new(connection.reader), MOST, false);
155 .next()
156 .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "No address for the relay"))?;
157 let tcp = TcpStream::connect_timeout(&target, CONNECT)?;
158 tcp.set_nodelay(true)?;
159 tcp.set_read_timeout(Some(QUIET))?;
160 tcp.set_write_timeout(Some(QUIET))?;
161 let (reading, mut writing): (Box<dyn Read + Send>, Box<dyn Write + Send>) = if address.tls {
162 let name = ServerName::try_from(address.host.clone())
163 .map_err(|error| io::Error::new(io::ErrorKind::InvalidInput, error))?;
164 let mut connection = ClientConnection::new(tls(), name).map_err(io::Error::other)?;
165 while connection.is_handshaking() {
166 connection.complete_io(&mut &tcp)?;
167 }
168 let connection = Arc::new(Mutex::new(connection));
169 (
170 Box::new(TlsReader {
171 tcp: tcp.try_clone()?,
172 tls: Arc::clone(&connection),
173 plain: Vec::new(),
174 at: 0,
175 }),
176 Box::new(TlsWriter {
177 tcp: tcp.try_clone()?,
178 tls: connection,
179 }),
180 )
181 } else {
182 (Box::new(tcp.try_clone()?), Box::new(tcp.try_clone()?))
183 };
184 let mut key = [0; 16];
185 getrandom::fill(&mut key).map_err(|_| io::Error::other("System random source failed"))?;
186 let key = base64::engine::general_purpose::STANDARD.encode(key);
187 write!(
188 writing,
189 "GET {path} HTTP/1.1\r\nHost: {}\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\
190 Sec-WebSocket-Key: {key}\r\nSec-WebSocket-Version: 13\r\nUser-Agent: Snowbound/{}\r\n\r\n",
191 address.authority,
192 env!("CARGO_PKG_VERSION"),
193 )?;
194 let mut reader = BufReader::new(reading);
195 let head = ws::head(&mut reader)?;
196 let status = head
197 .split(' ')
198 .nth(1)
199 .and_then(|status| status.parse().ok())
200 .unwrap_or(0);
201 if status != 101 {
202 let retry = ws::header(&head, "Retry-After")
203 .and_then(|seconds| seconds.parse().ok())
204 .map(Duration::from_secs);
205 return Err(Failure::Refused(status, retry));
206 }
207 if ws::header(&head, "Sec-WebSocket-Accept") != Some(ws::accept(&key).as_str()) {
208 return Err(io::Error::new(io::ErrorKind::InvalidData, "Not a relay").into());
209 }
210 let socket = Arc::new(Socket {105 let socket = Arc::new(Socket {
211 send: Mutex::new(writing),106 send: Mutex::new(connection.writer),
212 tcp,107 close: connection.close,
213 owner,108 owner,
214 links: Mutex::default(),109 links: Mutex::default(),
215 });110 });
216 Ok((socket, ws::Reader::new(reader, MOST, false)))111 Ok((socket, reader))
217}
218
219/// The certificate authorities the system trusts, then Mozilla's for a system whose store is
220/// missing or stale, as updates trust them.
221fn tls() -> Arc<ClientConfig> {
222 static CONFIG: OnceLock<Arc<ClientConfig>> = OnceLock::new();
223 Arc::clone(CONFIG.get_or_init(|| {
224 let mut roots = RootCertStore::empty();
225 roots.add_parsable_certificates(rustls_native_certs::load_native_certs().certs);
226 roots.add_parsable_certificates(webpki_root_certs::TLS_SERVER_ROOT_CERTS.iter().cloned());
227 let provider = Arc::new(rustls::crypto::ring::default_provider());
228 Arc::new(
229 ClientConfig::builder_with_provider(provider)
230 .with_safe_default_protocol_versions()
231 .expect("ring speaks TLS 1.2 and 1.3")
232 .with_root_certificates(roots)
233 .with_no_client_auth(),
234 )
235 }))
236}
237
238/// TLS read on one thread while another writes: the socket is read without the lock, and
239/// what arrives is decrypted under it.
240struct TlsReader {
241 tcp: TcpStream,
242 tls: Arc<Mutex<ClientConnection>>,
243 plain: Vec<u8>,
244 at: usize,
245}
246
247impl Read for TlsReader {
248 fn read(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
249 while self.at == self.plain.len() {
250 self.plain.clear();
251 self.at = 0;
252 let mut raw = vec![0; 16 << 10];
253 let length = self.tcp.read(&mut raw)?;
254 if length == 0 {
255 return Ok(0);
256 }
257 let mut tls = self.tls.lock().unwrap();
258 let mut arrived = &raw[..length];
259 let mut closed = false;
260 while !arrived.is_empty() {
261 tls.read_tls(&mut arrived)?;
262 tls.process_new_packets()
263 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
264 let mut chunk = [0; 4096];
265 loop {
266 match tls.reader().read(&mut chunk) {
267 Ok(0) => {
268 closed = true;
269 break;
270 }
271 Ok(length) => self.plain.extend_from_slice(&chunk[..length]),
272 Err(error) if error.kind() == io::ErrorKind::WouldBlock => break,
273 Err(error) => return Err(error),
274 }
275 }
276 }
277 while tls.wants_write() {
278 tls.write_tls(&mut &self.tcp)?;
279 }
280 if closed && self.plain.is_empty() {
281 return Ok(0);
282 }
283 }
284 let length = buffer.len().min(self.plain.len() - self.at);
285 buffer[..length].copy_from_slice(&self.plain[self.at..self.at + length]);
286 self.at += length;
287 Ok(length)
288 }
289}
290
291struct TlsWriter {
292 tcp: TcpStream,
293 tls: Arc<Mutex<ClientConnection>>,
294}
295
296impl Write for TlsWriter {
297 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
298 let mut tls = self.tls.lock().unwrap();
299 let length = tls.writer().write(bytes)?;
300 while tls.wants_write() {
301 tls.write_tls(&mut &self.tcp)?;
302 }
303 Ok(length)
304 }
305
306 fn flush(&mut self) -> io::Result<()> {
307 Ok(())
308 }
309}112}
310113
311/// One connection to a relay's room.114/// One connection to a relay's room.
312pub(super) struct Socket {115pub(super) struct Socket {
313 send: Mutex<Box<dyn Write + Send>>,116 send: Mutex<Box<dyn Write + Send>>,
314 tcp: TcpStream,117 close: Box<dyn Fn() + Send + Sync>,
315 /// Whether this end claimed the room for its code, and so tells the relay who knew it.118 /// Whether this end claimed the room for its code, and so tells the relay who knew it.
316 owner: bool,119 owner: bool,
317 links: Mutex<Links>,120 links: Mutex<Links>,
...@@ -336,7 +139,7 @@ impl Socket {...@@ -336,7 +139,7 @@ impl Socket {
336 }139 }
337140
338 pub(super) fn hang_up(&self) {141 pub(super) fn hang_up(&self) {
339 let _ = self.tcp.shutdown(Shutdown::Both);142 (self.close)();
340 }143 }
341144
342 fn forget(&self, slot: u32) {145 fn forget(&self, slot: u32) {
crates/notebook/src/live/share.rs+5-5
...@@ -114,8 +114,8 @@ pub enum Refusal {...@@ -114,8 +114,8 @@ pub enum Refusal {
114 TooMany(Option<Duration>),114 TooMany(Option<Duration>),
115 /// The relay is full.115 /// The relay is full.
116 Busy,116 Busy,
117 /// The relay couldn't be reached, and no one answered on this network.117 /// The relay couldn't be reached, and why, and no one answered on this network.
118 Unreachable,118 Unreachable(super::Trouble),
119 /// The relay let this end in, but no one answered.119 /// The relay let this end in, but no one answered.
120 TimedOut,120 TimedOut,
121}121}
...@@ -152,7 +152,7 @@ pub fn join(...@@ -152,7 +152,7 @@ pub fn join(
152 }152 }
153 },153 },
154 )154 )
155 .map_err(|_| Refusal::Unreachable)?;155 .map_err(|_| Refusal::Unreachable(super::Trouble::Other))?;
156 let start = Instant::now();156 let start = Instant::now();
157 loop {157 loop {
158 if let Ok(welcome) = welcome.try_recv() {158 if let Ok(welcome) = welcome.try_recv() {
...@@ -174,8 +174,8 @@ pub fn join(...@@ -174,8 +174,8 @@ pub fn join(
174 Relayed::Refused(410, _) => return Err(Refusal::Expired),174 Relayed::Refused(410, _) => return Err(Refusal::Expired),
175 Relayed::Refused(429, wait) => return Err(Refusal::TooMany(wait)),175 Relayed::Refused(429, wait) => return Err(Refusal::TooMany(wait)),
176 Relayed::Refused(503, _) if settled => return Err(Refusal::Busy),176 Relayed::Refused(503, _) if settled => return Err(Refusal::Busy),
177 Relayed::Unreachable if waited > Duration::from_secs(10) => {177 Relayed::Unreachable(trouble) if waited > Duration::from_secs(10) => {
178 return Err(Refusal::Unreachable);178 return Err(Refusal::Unreachable(trouble));
179 }179 }
180 Relayed::Unknown if relay.is_none() && waited > Duration::from_secs(10) => {180 Relayed::Unknown if relay.is_none() && waited > Duration::from_secs(10) => {
181 return Err(Refusal::NoOne);181 return Err(Refusal::NoOne);
crates/notebook/src/live/transport.rs created+628
...@@ -0,0 +1,628 @@
1//! The connection to a relay, wherever the network lets one through: straight to it, or
2//! through the proxy `proxy` names by `CONNECT`, over TLS that trusts what the system trusts
3//! (so a proxy that inspects HTTPS with its own authority works once the system trusts it);
4//! as a WebSocket, or where something on the way refuses WebSockets, as HTTPS requests the
5//! relay answers the same messages over (`GET` waits for what is to come, `POST` sends). What
6//! fails is told apart, so a person can be told what to try.
7
8use super::proxy::{self, Proxy};
9use ::relay::ws;
10use base64::Engine;
11use rustls::{ClientConfig, ClientConnection, RootCertStore, pki_types::ServerName};
12use std::{
13 io::{self, BufReader, Read, Write},
14 net::{Shutdown, TcpStream, ToSocketAddrs},
15 sync::{
16 Arc, Condvar, Mutex, OnceLock,
17 atomic::{AtomicBool, Ordering},
18 },
19 thread,
20 time::Duration,
21};
22
23const CONNECT: Duration = Duration::from_secs(10);
24/// How long a connection may hear nothing: the relay hears a ping every 30 s and answers.
25pub(super) const QUIET: Duration = Duration::from_secs(75);
26/// The most a poll's answer holds.
27const MOST: usize = 4 << 20;
28/// Set once WebSockets were refused, so later connections go straight to polling.
29static POLLING: AtomicBool = AtomicBool::new(false);
30
31/// Why a relay could not be reached, as a person can act on it.
32#[derive(Clone, Copy, Debug, PartialEq, Eq)]
33pub enum Trouble {
34 /// The relay's name, or the proxy's, didn't resolve.
35 Dns,
36 /// Nothing answered at the relay's address: a firewall, or the relay is down.
37 Unreachable,
38 TimedOut,
39 /// The proxy itself couldn't be reached.
40 ProxyUnreachable,
41 /// The proxy asks for a name and password (407).
42 ProxyAuthentication,
43 /// The proxy refused to connect to the relay, with this status.
44 ProxyRefused(u16),
45 /// The relay's certificate wasn't one the system trusts: something on the way presents its
46 /// own.
47 Certificate,
48 /// Something on the way refused both WebSockets and plain requests to the relay.
49 Blocked,
50 Other,
51}
52
53pub(super) enum Failure {
54 Trouble(Trouble, io::Error),
55 /// The relay answered, with this HTTP status and how long to wait.
56 Refused(u16, Option<Duration>),
57}
58
59impl Failure {
60 fn other(error: io::Error) -> Self {
61 Failure::Trouble(Trouble::Other, error)
62 }
63}
64
65impl From<io::Error> for Failure {
66 fn from(error: io::Error) -> Self {
67 Failure::other(error)
68 }
69}
70
71/// A relay's address: `ws://` or `wss://`, a host, and the path the relay's `/v1/` follows.
72#[derive(Debug, PartialEq)]
73pub(super) struct Address {
74 pub tls: bool,
75 /// The host and port as the URL gave them, for the `Host` header.
76 pub authority: String,
77 pub host: String,
78 pub port: u16,
79 pub path: String,
80}
81
82pub(super) fn parse(url: &str) -> io::Result<Address> {
83 let bad = || io::Error::new(io::ErrorKind::InvalidInput, format!("Not a relay: {url}"));
84 let (tls, rest) = match url.split_once("://") {
85 Some(("wss", rest)) => (true, rest),
86 Some(("ws", rest)) => (false, rest),
87 _ => return Err(bad()),
88 };
89 let (authority, path) = rest.split_at(rest.find('/').unwrap_or(rest.len()));
90 let (host, port) = match authority.rsplit_once(':') {
91 Some((host, port)) if !port.contains(']') => (host, port.parse().map_err(|_| bad())?),
92 _ => (authority, if tls { 443 } else { 80 }),
93 };
94 let host = host.trim_start_matches('[').trim_end_matches(']');
95 if host.is_empty() {
96 return Err(bad());
97 }
98 Ok(Address {
99 tls,
100 authority: authority.into(),
101 host: host.into(),
102 port,
103 path: path.trim_end_matches('/').into(),
104 })
105}
106
107/// A connection to the relay: what it writes the relay's messages into, what it reads them
108/// from, and how to hang up.
109pub(super) struct Connection {
110 pub writer: Box<dyn Write + Send>,
111 pub reader: Box<dyn Read + Send>,
112 pub close: Box<dyn Fn() + Send + Sync>,
113}
114
115/// Opens the relay's `path` at `address`, as a WebSocket, else as HTTPS requests where
116/// something on the way refused the WebSocket.
117pub(super) fn connect(address: &Address, path: &str) -> Result<Connection, Failure> {
118 if !POLLING.load(Ordering::Acquire) {
119 match websocket(address, path) {
120 Err(Failure::Refused(status, _)) if !relay_status(status) => {
121 POLLING.store(true, Ordering::Release);
122 }
123 Err(Failure::Trouble(Trouble::Other, _)) => {
124 POLLING.store(true, Ordering::Release);
125 }
126 connected => return connected,
127 }
128 }
129 poll(address, path).map_err(|failure| match failure {
130 Failure::Refused(status, _) if !relay_status(status) => Failure::Trouble(
131 Trouble::Blocked,
132 io::Error::other(format!("Refused with {status}")),
133 ),
134 failure => failure,
135 })
136}
137
138/// Whether the relay itself answers with `status`, rather than something on the way.
139fn relay_status(status: u16) -> bool {
140 matches!(status, 404 | 410 | 429 | 503)
141}
142
143/// A stream to the relay, through the proxy where there is one, over TLS where its address
144/// asks.
145struct Stream {
146 reader: Box<dyn Read + Send>,
147 writer: Box<dyn Write + Send>,
148 tcp: TcpStream,
149}
150
151fn open(address: &Address) -> Result<Stream, Failure> {
152 let proxy = proxy::for_host(&address.host, address.tls);
153 let (host, port) = match &proxy {
154 Some(proxy) => (proxy.host.as_str(), proxy.port),
155 None => (address.host.as_str(), address.port),
156 };
157 let reached = |trouble| match proxy {
158 Some(_) => Trouble::ProxyUnreachable,
159 None => trouble,
160 };
161 let target = (host, port)
162 .to_socket_addrs()
163 .map_err(|error| Failure::Trouble(reached(Trouble::Dns), error))?
164 .next()
165 .ok_or_else(|| Failure::Trouble(reached(Trouble::Dns), io::ErrorKind::NotFound.into()))?;
166 let tcp = TcpStream::connect_timeout(&target, CONNECT).map_err(|error| {
167 let trouble = match error.kind() {
168 io::ErrorKind::TimedOut => Trouble::TimedOut,
169 _ => Trouble::Unreachable,
170 };
171 Failure::Trouble(reached(trouble), error)
172 })?;
173 tcp.set_nodelay(true)?;
174 tcp.set_read_timeout(Some(QUIET))?;
175 tcp.set_write_timeout(Some(QUIET))?;
176 if let Some(proxy) = &proxy {
177 tunnel(&tcp, proxy, address)?;
178 }
179 if !address.tls {
180 return Ok(Stream {
181 reader: Box::new(tcp.try_clone()?),
182 writer: Box::new(tcp.try_clone()?),
183 tcp,
184 });
185 }
186 let name = ServerName::try_from(address.host.clone())
187 .map_err(|error| io::Error::new(io::ErrorKind::InvalidInput, error))?;
188 let mut connection = ClientConnection::new(tls(), name).map_err(io::Error::other)?;
189 while connection.is_handshaking() {
190 connection.complete_io(&mut &tcp).map_err(|error| {
191 let rejected = error
192 .get_ref()
193 .and_then(|inner| inner.downcast_ref::<rustls::Error>())
194 .is_some_and(|tls| matches!(tls, rustls::Error::InvalidCertificate(_)));
195 match rejected {
196 true => Failure::Trouble(Trouble::Certificate, error),
197 false => Failure::other(error),
198 }
199 })?;
200 }
201 let connection = Arc::new(Mutex::new(connection));
202 Ok(Stream {
203 reader: Box::new(TlsReader {
204 tcp: tcp.try_clone()?,
205 tls: Arc::clone(&connection),
206 plain: Vec::new(),
207 at: 0,
208 }),
209 writer: Box::new(TlsWriter {
210 tcp: tcp.try_clone()?,
211 tls: connection,
212 }),
213 tcp,
214 })
215}
216
217/// Asks `proxy` on `tcp` to connect through to the relay.
218fn tunnel(tcp: &TcpStream, proxy: &Proxy, address: &Address) -> Result<(), Failure> {
219 let target = format!("{}:{}", address.host, address.port);
220 let mut request = format!("CONNECT {target} HTTP/1.1\r\nHost: {target}\r\n");
221 if let Some(authorization) = proxy.authorization() {
222 request += &format!("Proxy-Authorization: {authorization}\r\n");
223 }
224 request += "\r\n";
225 (&mut &*tcp)
226 .write_all(request.as_bytes())
227 .map_err(|error| Failure::Trouble(Trouble::ProxyUnreachable, error))?;
228 let head =
229 ws::head(&mut &*tcp).map_err(|error| Failure::Trouble(Trouble::ProxyUnreachable, error))?;
230 match status(&head) {
231 200 => Ok(()),
232 407 => Err(Failure::Trouble(
233 Trouble::ProxyAuthentication,
234 io::Error::new(io::ErrorKind::PermissionDenied, "The proxy asks for a name"),
235 )),
236 status => Err(Failure::Trouble(
237 Trouble::ProxyRefused(status),
238 io::Error::other(format!("The proxy answered {status}")),
239 )),
240 }
241}
242
243fn status(head: &str) -> u16 {
244 head.split(' ')
245 .nth(1)
246 .and_then(|status| status.parse().ok())
247 .unwrap_or(0)
248}
249
250fn retry(head: &str) -> Option<Duration> {
251 ws::header(head, "Retry-After")
252 .and_then(|seconds| seconds.parse().ok())
253 .map(Duration::from_secs)
254}
255
256/// Opens a WebSocket to `path` at `address`.
257fn websocket(address: &Address, path: &str) -> Result<Connection, Failure> {
258 let Stream {
259 reader,
260 mut writer,
261 tcp,
262 } = open(address)?;
263 let mut key = [0; 16];
264 getrandom::fill(&mut key).map_err(|_| io::Error::other("System random source failed"))?;
265 let key = base64::engine::general_purpose::STANDARD.encode(key);
266 write!(
267 writer,
268 "GET {path} HTTP/1.1\r\nHost: {}\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\
269 Sec-WebSocket-Key: {key}\r\nSec-WebSocket-Version: 13\r\nUser-Agent: Snowbound/{}\r\n\r\n",
270 address.authority,
271 env!("CARGO_PKG_VERSION"),
272 )?;
273 let mut reader = BufReader::new(reader);
274 let head = ws::head(&mut reader)?;
275 if status(&head) != 101 {
276 return Err(Failure::Refused(status(&head), retry(&head)));
277 }
278 if ws::header(&head, "Sec-WebSocket-Accept") != Some(ws::accept(&key).as_str()) {
279 return Err(io::Error::new(io::ErrorKind::InvalidData, "Not a relay").into());
280 }
281 Ok(Connection {
282 writer,
283 reader: Box::new(reader),
284 close: Box::new(move || {
285 let _ = tcp.shutdown(Shutdown::Both);
286 }),
287 })
288}
289
290/// An HTTP/1.1 connection that answers one request after another.
291struct Http {
292 reader: BufReader<Box<dyn Read + Send>>,
293 writer: Box<dyn Write + Send>,
294 tcp: TcpStream,
295}
296
297impl Http {
298 fn open(address: &Address) -> Result<Self, Failure> {
299 let Stream {
300 reader,
301 writer,
302 tcp,
303 } = open(address)?;
304 Ok(Self {
305 reader: BufReader::new(reader),
306 writer,
307 tcp,
308 })
309 }
310
311 /// Sends `method` on `path` with `body`: the status, the head and the body answered.
312 fn ask(
313 &mut self,
314 address: &Address,
315 method: &str,
316 path: &str,
317 body: &[u8],
318 ) -> io::Result<(u16, String, Vec<u8>)> {
319 write!(
320 self.writer,
321 "{method} {path} HTTP/1.1\r\nHost: {}\r\nUser-Agent: Snowbound/{}\r\n\
322 Content-Length: {}\r\nContent-Type: application/octet-stream\r\n\r\n",
323 address.authority,
324 env!("CARGO_PKG_VERSION"),
325 body.len()
326 )?;
327 self.writer.write_all(body)?;
328 let head = ws::head(&mut self.reader)?;
329 let length: usize = ws::header(&head, "Content-Length")
330 .and_then(|length| length.parse().ok())
331 .unwrap_or(0);
332 if length > MOST {
333 return Err(io::Error::new(
334 io::ErrorKind::InvalidData,
335 "An answer too large",
336 ));
337 }
338 let mut answer = vec![0; length];
339 self.reader.read_exact(&mut answer)?;
340 Ok((status(&head), head, answer))
341 }
342}
343
344/// What a polled session shares between its reading, its sending and hanging up.
345struct Session {
346 /// Messages written and not yet sent.
347 pending: Mutex<Vec<u8>>,
348 ready: Condvar,
349 closed: AtomicBool,
350 /// The connections open now, to shut on hanging up.
351 open: Mutex<Vec<TcpStream>>,
352}
353
354impl Session {
355 fn close(&self) {
356 self.closed.store(true, Ordering::Release);
357 self.ready.notify_all();
358 for tcp in self.open.lock().unwrap().drain(..) {
359 let _ = tcp.shutdown(Shutdown::Both);
360 }
361 }
362}
363
364/// Joins `path` as HTTPS requests: the relay's answer names a session, whose messages a
365/// `GET` waits for and a `POST` sends, each a run of WebSocket frames as on a WebSocket.
366fn poll(address: &Address, path: &str) -> Result<Connection, Failure> {
367 let mut http = Http::open(address)?;
368 let joined = match path.contains('?') {
369 true => format!("{path}&poll=1"),
370 false => format!("{path}?poll=1"),
371 };
372 let (status, head, body) = http.ask(address, "GET", &joined, &[])?;
373 if status != 200 {
374 return Err(Failure::Refused(status, retry(&head)));
375 }
376 let session_id = String::from_utf8_lossy(&body)
377 .strip_prefix("session ")
378 .map(|id| id.trim().to_owned())
379 .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "Not a relay"))?;
380 let at = format!("{}/v1/poll/{session_id}", address.path);
381 let session = Arc::new(Session {
382 pending: Mutex::default(),
383 ready: Condvar::new(),
384 closed: AtomicBool::new(false),
385 open: Mutex::new(vec![http.tcp.try_clone()?]),
386 });
387 let address = Arc::new(Address {
388 tls: address.tls,
389 authority: address.authority.clone(),
390 host: address.host.clone(),
391 port: address.port,
392 path: address.path.clone(),
393 });
394 // Sends what is written, a batch a request, on a connection of its own.
395 let (sending, to, posting) = (Arc::clone(&session), at.clone(), Arc::clone(&address));
396 thread::Builder::new()
397 .name("live relay post".into())
398 .spawn(move || post(&sending, &posting, &to))?;
399 let closing = Arc::clone(&session);
400 Ok(Connection {
401 writer: Box::new(PollWriter(Arc::clone(&session))),
402 reader: Box::new(PollReader {
403 session,
404 address,
405 at,
406 http: Some(http),
407 arrived: Vec::new(),
408 read: 0,
409 }),
410 close: Box::new(move || closing.close()),
411 })
412}
413
414fn post(session: &Session, address: &Address, at: &str) {
415 let mut http: Option<Http> = None;
416 loop {
417 let batch = {
418 let mut pending = session.pending.lock().unwrap();
419 while pending.is_empty() && !session.closed.load(Ordering::Acquire) {
420 pending = session.ready.wait(pending).unwrap();
421 }
422 if session.closed.load(Ordering::Acquire) {
423 return;
424 }
425 std::mem::take(&mut *pending)
426 };
427 // A connection the relay or the proxy closed meanwhile is opened again, once.
428 let sent = (0..2).any(|_| {
429 if http.is_none() {
430 http = Http::open(address).ok();
431 if let (Some(opened), Ok(mut open)) = (&http, session.open.lock()) {
432 open.extend(opened.tcp.try_clone());
433 }
434 }
435 let answered = http
436 .as_mut()
437 .map(|connection| connection.ask(address, "POST", at, &batch));
438 match answered {
439 Some(Ok((200 | 204, ..))) => true,
440 Some(Ok((410, ..))) => {
441 session.close();
442 true
443 }
444 _ => {
445 http = None;
446 false
447 }
448 }
449 });
450 if !sent {
451 session.close();
452 return;
453 }
454 }
455}
456
457struct PollWriter(Arc<Session>);
458
459impl Write for PollWriter {
460 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
461 if self.0.closed.load(Ordering::Acquire) {
462 return Err(io::ErrorKind::BrokenPipe.into());
463 }
464 self.0.pending.lock().unwrap().extend_from_slice(bytes);
465 self.0.ready.notify_all();
466 Ok(bytes.len())
467 }
468
469 fn flush(&mut self) -> io::Result<()> {
470 Ok(())
471 }
472}
473
474/// Reads what each `GET` brings, waiting for the next when it is all read.
475struct PollReader {
476 session: Arc<Session>,
477 address: Arc<Address>,
478 at: String,
479 http: Option<Http>,
480 arrived: Vec<u8>,
481 read: usize,
482}
483
484impl Read for PollReader {
485 fn read(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
486 while self.read == self.arrived.len() {
487 if self.session.closed.load(Ordering::Acquire) {
488 return Ok(0);
489 }
490 if self.http.is_none() {
491 let opened = Http::open(&self.address).map_err(|failure| match failure {
492 Failure::Trouble(_, error) => error,
493 Failure::Refused(status, _) => io::Error::other(format!("{status}")),
494 })?;
495 self.session
496 .open
497 .lock()
498 .unwrap()
499 .extend(opened.tcp.try_clone());
500 self.http = Some(opened);
501 }
502 let answered =
503 self.http
504 .as_mut()
505 .expect("opened above")
506 .ask(&self.address, "GET", &self.at, &[]);
507 match answered {
508 Ok((200, _, body)) => {
509 self.arrived = body;
510 self.read = 0;
511 }
512 Ok((410, ..)) => {
513 self.session.close();
514 return Ok(0);
515 }
516 Ok((status, ..)) => {
517 return Err(io::Error::other(format!("The relay answered {status}")));
518 }
519 Err(error) => {
520 // Opened again once; a second failure ends the session.
521 if self.http.take().is_none() {
522 return Err(error);
523 }
524 self.http = Http::open(&self.address).ok();
525 if self.http.is_none() {
526 return Err(error);
527 }
528 }
529 }
530 }
531 let length = buffer.len().min(self.arrived.len() - self.read);
532 buffer[..length].copy_from_slice(&self.arrived[self.read..self.read + length]);
533 self.read += length;
534 Ok(length)
535 }
536}
537
538/// The certificate authorities the system trusts, then Mozilla's for a system whose store is
539/// missing or stale, as updates trust them.
540fn tls() -> Arc<ClientConfig> {
541 static CONFIG: OnceLock<Arc<ClientConfig>> = OnceLock::new();
542 Arc::clone(CONFIG.get_or_init(|| {
543 let mut roots = RootCertStore::empty();
544 roots.add_parsable_certificates(rustls_native_certs::load_native_certs().certs);
545 roots.add_parsable_certificates(webpki_root_certs::TLS_SERVER_ROOT_CERTS.iter().cloned());
546 let provider = Arc::new(rustls::crypto::ring::default_provider());
547 Arc::new(
548 ClientConfig::builder_with_provider(provider)
549 .with_safe_default_protocol_versions()
550 .expect("ring speaks TLS 1.2 and 1.3")
551 .with_root_certificates(roots)
552 .with_no_client_auth(),
553 )
554 }))
555}
556
557/// TLS read on one thread while another writes: the socket is read without the lock, and
558/// what arrives is decrypted under it.
559struct TlsReader {
560 tcp: TcpStream,
561 tls: Arc<Mutex<ClientConnection>>,
562 plain: Vec<u8>,
563 at: usize,
564}
565
566impl Read for TlsReader {
567 fn read(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
568 while self.at == self.plain.len() {
569 self.plain.clear();
570 self.at = 0;
571 let mut raw = vec![0; 16 << 10];
572 let length = self.tcp.read(&mut raw)?;
573 if length == 0 {
574 return Ok(0);
575 }
576 let mut tls = self.tls.lock().unwrap();
577 let mut arrived = &raw[..length];
578 let mut closed = false;
579 while !arrived.is_empty() {
580 tls.read_tls(&mut arrived)?;
581 tls.process_new_packets()
582 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
583 let mut chunk = [0; 4096];
584 loop {
585 match tls.reader().read(&mut chunk) {
586 Ok(0) => {
587 closed = true;
588 break;
589 }
590 Ok(length) => self.plain.extend_from_slice(&chunk[..length]),
591 Err(error) if error.kind() == io::ErrorKind::WouldBlock => break,
592 Err(error) => return Err(error),
593 }
594 }
595 }
596 while tls.wants_write() {
597 tls.write_tls(&mut &self.tcp)?;
598 }
599 if closed && self.plain.is_empty() {
600 return Ok(0);
601 }
602 }
603 let length = buffer.len().min(self.plain.len() - self.at);
604 buffer[..length].copy_from_slice(&self.plain[self.at..self.at + length]);
605 self.at += length;
606 Ok(length)
607 }
608}
609
610struct TlsWriter {
611 tcp: TcpStream,
612 tls: Arc<Mutex<ClientConnection>>,
613}
614
615impl Write for TlsWriter {
616 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
617 let mut tls = self.tls.lock().unwrap();
618 let length = tls.writer().write(bytes)?;
619 while tls.wants_write() {
620 tls.write_tls(&mut &self.tcp)?;
621 }
622 Ok(length)
623 }
624
625 fn flush(&mut self) -> io::Result<()> {
626 Ok(())
627 }
628}
crates/notebook/tests/live_networks.rs created+168
...@@ -0,0 +1,168 @@
1//! Live Share on hostile networks: through a proxy that asks for a password, through one that
2//! refuses WebSockets (the relay is then reached by plain requests), and the troubles named
3//! when the relay can't be reached. One test, as the proxy chosen holds for the process.
4#![cfg(feature = "live")]
5
6use notebook::live::{
7 Trouble,
8 proxy::{self, Proxy},
9 share::{self, Refusal, Sharing},
10};
11use std::{
12 io::{self, Read, Write},
13 net::{SocketAddr, TcpListener, TcpStream},
14 sync::{
15 Arc,
16 atomic::{AtomicUsize, Ordering},
17 },
18 thread,
19};
20
21#[path = "support/live.rs"]
22mod live;
23use live::*;
24
25/// What a proxy on this computer did.
26#[derive(Default)]
27struct Seen {
28 tunnels: AtomicUsize,
29 refused: AtomicUsize,
30}
31
32/// An HTTP proxy that tunnels `CONNECT`s, asking for `credentials` where given and, with
33/// `block_websockets`, refusing a WebSocket's upgrade inside the tunnel, as a proxy that
34/// inspects HTTPS does.
35fn proxy(credentials: Option<&'static str>, block_websockets: bool) -> (SocketAddr, Arc<Seen>) {
36 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
37 let address = listener.local_addr().unwrap();
38 let seen = Arc::new(Seen::default());
39 let counting = Arc::clone(&seen);
40 thread::spawn(move || {
41 for client in listener.incoming().flatten() {
42 let seen = Arc::clone(&counting);
43 thread::spawn(move || {
44 let _ = tunnel(client, credentials, block_websockets, &seen);
45 });
46 }
47 });
48 (address, seen)
49}
50
51fn tunnel(
52 mut client: TcpStream,
53 credentials: Option<&str>,
54 block_websockets: bool,
55 seen: &Seen,
56) -> io::Result<()> {
57 let head = relay::ws::head(&mut client)?;
58 let target = head.split(' ').nth(1).unwrap_or_default().to_owned();
59 let authorized = credentials
60 .is_none_or(|expected| relay::ws::header(&head, "Proxy-Authorization") == Some(expected));
61 if !head.starts_with("CONNECT ") || !authorized {
62 client.write_all(
63 b"HTTP/1.1 407 Proxy Authentication Required\r\nProxy-Authenticate: Basic\r\n\
64 Content-Length: 0\r\n\r\n",
65 )?;
66 return Ok(());
67 }
68 let mut upstream = TcpStream::connect(&target)?;
69 client.write_all(b"HTTP/1.1 200 Connection established\r\n\r\n")?;
70 seen.tunnels.fetch_add(1, Ordering::Relaxed);
71 if block_websockets {
72 let request = relay::ws::head(&mut client)?;
73 if relay::ws::header(&request, "Upgrade").is_some() {
74 seen.refused.fetch_add(1, Ordering::Relaxed);
75 client.write_all(b"HTTP/1.1 403 Forbidden\r\nContent-Length: 0\r\n\r\n")?;
76 return Ok(());
77 }
78 upstream.write_all(request.as_bytes())?;
79 }
80 let (mut from, mut to) = (client.try_clone()?, upstream.try_clone()?);
81 thread::spawn(move || {
82 let _ = io::copy(&mut from, &mut to);
83 let _ = to.shutdown(std::net::Shutdown::Both);
84 });
85 let mut buffer = [0; 16 << 10];
86 loop {
87 let length = upstream.read(&mut buffer)?;
88 if length == 0 {
89 return Ok(());
90 }
91 client.write_all(&buffer[..length])?;
92 }
93}
94
95fn through(address: SocketAddr, credentials: Option<(&str, &str)>) {
96 proxy::use_proxy(Some(Some(Proxy {
97 host: address.ip().to_string(),
98 port: address.port(),
99 credentials: credentials.map(|(name, password)| (name.into(), password.into())),
100 })));
101}
102
103/// Shares a notebook and has a guest join it and publish an edit, all through `url`.
104fn share_and_edit(directory: &std::path::Path, url: &str) {
105 let folder = notebook(directory);
106 let host = host(
107 &folder,
108 &directory.join("host"),
109 &Sharing::new("").unwrap(),
110 url,
111 );
112 let (guest, notebook) = guest("Grace", &code(&host), url, &directory.join("grace"));
113 let section = open(&notebook, &guest, "Garden.one", None);
114 let file = folder.join("Garden.one");
115 let id = replace(&section, &std::fs::read(&file).unwrap(), 0..8, "Through");
116 published(&section, id);
117 assert_eq!(
118 server::text(&std::fs::read(&file).unwrap()).2,
119 "Through text"
120 );
121}
122
123fn refusal(code: &str, url: &str) -> Refusal {
124 share::join(hello("Mallory"), code, "", None, Some(url)).unwrap_err()
125}
126
127#[test]
128fn hostile_networks() {
129 let directory = tempfile::tempdir().unwrap();
130 let url = relay(Default::default());
131 let code = notebook::live::code::format(412, "4MZ9XR").unwrap();
132
133 // A proxy that asks for a password: refused without, through with it.
134 let (address, seen) = proxy(Some("Basic YWRhOnNlY3JldA=="), false);
135 through(address, None);
136 assert_eq!(
137 refusal(&code, &url),
138 Refusal::Unreachable(Trouble::ProxyAuthentication)
139 );
140 through(address, Some(("ada", "secret")));
141 share_and_edit(&directory.path().join("connect"), &url);
142 assert!(
143 seen.tunnels.load(Ordering::Relaxed) >= 3,
144 "host, guest and joiner tunnel"
145 );
146
147 // A proxy that refuses WebSockets: the relay is reached by requests instead.
148 let (address, seen) = proxy(None, true);
149 through(address, None);
150 share_and_edit(&directory.path().join("blocked"), &url);
151 assert!(seen.refused.load(Ordering::Relaxed) >= 1);
152
153 // A proxy that isn't there, and a relay whose name doesn't resolve.
154 let closed = TcpListener::bind("127.0.0.1:0")
155 .unwrap()
156 .local_addr()
157 .unwrap();
158 through(closed, None);
159 assert_eq!(
160 refusal(&code, &url),
161 Refusal::Unreachable(Trouble::ProxyUnreachable)
162 );
163 proxy::use_proxy(Some(None));
164 assert_eq!(
165 refusal(&code, "wss://relay.invalid"),
166 Refusal::Unreachable(Trouble::Dns)
167 );
168}
crates/notebook/tests/live_share.rs+7-151
...@@ -4,159 +4,15 @@...@@ -4,159 +4,15 @@
4#![cfg(feature = "live")]4#![cfg(feature = "live")]
55
6use notebook::{6use notebook::{
7 EditStatus, Replica,7 live::share::{self, Refusal, Sharing},
8 live::{8 session::{Notebook, SyncState},
9 Hello,
10 share::{self, Guest, Host, Refusal, Sharing},
11 },
12 session::{Notebook, Section, SyncState},
13};9};
14use onestore::{10use onestore::Arena;
15 Arena, ExGuid,11use std::sync::Arc;
16 op::{Edit, Op, PageOp},
17 protected::{Key, rekey},
18};
19use std::{
20 net::TcpListener,
21 path::Path,
22 sync::Arc,
23 thread,
24 time::{Duration, Instant},
25};
26
27#[path = "support/server.rs"]
28mod server;
29
30const PASSWORD: &str = "fixture password";
31
32/// A relay on this computer with `config`'s limits: its URL.
33fn relay(config: relay::server::Config) -> String {
34 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
35 let url = format!("ws://{}", listener.local_addr().unwrap());
36 thread::spawn(move || relay::server::serve(listener, config));
37 url
38}
39
40/// Another secret than `secret`, as a guess makes one.
41fn mistaken(secret: &str) -> String {
42 let first = if secret.starts_with('A') { 'B' } else { 'A' };
43 format!("{first}{}", &secret[1..])
44}
45
46fn hello(name: &str) -> Hello {
47 Hello::new(name.into(), None).unwrap()
48}
49
50fn until(what: &str, done: impl Fn() -> bool) {
51 let deadline = Instant::now() + Duration::from_secs(30);
52 while !done() {
53 assert!(Instant::now() < deadline, "{what}");
54 thread::sleep(Duration::from_millis(20));
55 }
56}
5712
58/// A notebook folder holding `Garden.one`, a page reading "Original text", and a protected13#[path = "support/live.rs"]
59/// `Sealed.one` reading "Sealed text".14mod live;
60fn notebook(root: &Path) -> std::path::PathBuf {15use live::*;
61 let folder = root.join("Garden");
62 std::fs::create_dir_all(&folder).unwrap();
63 std::fs::write(
64 folder.join("Garden.one"),
65 onestore::create_section("Garden.one", "Original text", "Fixture").unwrap(),
66 )
67 .unwrap();
68 let plain = onestore::create_section("Sealed.one", "Sealed text", "Fixture").unwrap();
69 let key = Key::new(PASSWORD).unwrap();
70 std::fs::write(
71 folder.join("Sealed.one"),
72 rekey(&plain, None, Some(&key)).unwrap(),
73 )
74 .unwrap();
75 folder
76}
77
78fn host(folder: &Path, cache: &Path, sharing: &Sharing, url: &str) -> Host {
79 let storage = Notebook::open(folder, cache).unwrap().into_storage();
80 Host::start(
81 storage,
82 hello("Ada"),
83 sharing.clone(),
84 "Garden",
85 None,
86 Some(url),
87 || {},
88 )
89 .unwrap()
90}
91
92/// The host's code once the relay has numbered it.
93fn code(host: &Host) -> String {
94 until("the code was never numbered", || {
95 host.code().is_some_and(|code| share::code(&code).is_some())
96 });
97 host.code().unwrap()
98}
99
100/// `name` joins with `code` and opens the notebook in `cache`.
101fn guest(name: &str, code: &str, url: &str, cache: &Path) -> (Arc<Guest>, Notebook) {
102 let welcome = share::join(hello(name), code, "", None, Some(url)).unwrap();
103 assert_eq!(
104 (welcome.notebook.as_str(), welcome.host.as_str()),
105 ("Garden", "Ada")
106 );
107 let guest = Guest::start(
108 hello(name),
109 welcome.share,
110 welcome.secret,
111 None,
112 Some(url),
113 || {},
114 )
115 .unwrap();
116 until("the host was never met", || guest.host().is_some());
117 let notebook = Notebook::open_hosted(Arc::clone(&guest), cache).unwrap();
118 (guest, notebook)
119}
120
121fn open(notebook: &Notebook, guest: &Arc<Guest>, path: &str, key: Option<&Key>) -> Section {
122 let replica = notebook.replica_path(path).unwrap();
123 std::fs::create_dir_all(replica.parent().unwrap()).unwrap();
124 let replica = Replica::open_or_create(&replica, key, || notebook.read_section(path)).unwrap();
125 Section::resume_hosted(path.into(), replica, Arc::clone(guest), || {}).unwrap()
126}
127
128fn replace(section: &Section, image: &[u8], range: std::ops::Range<u32>, with: &str) -> u64 {
129 let (space, text, _) = server::text(image);
130 replaced(section, space, text, range, with)
131}
132
133fn replaced(
134 section: &Section,
135 space: ExGuid,
136 text: ExGuid,
137 range: std::ops::Range<u32>,
138 with: &str,
139) -> u64 {
140 let op = PageOp::Text {
141 text,
142 range,
143 with: with.into(),
144 };
145 let edit = Edit {
146 at: 134_000_000_000_000_000,
147 ops: vec![Op::Page { space, op }],
148 };
149 section.replica().apply("Guest", edit).unwrap()
150}
151
152fn published(section: &Section, id: u64) {
153 until("the edit was never published", || {
154 matches!(
155 section.status(id).unwrap(),
156 Some(EditStatus::Published { .. })
157 )
158 });
159}
16016
161/// A guest opens a section through the host, its edit lands in the host's file, and the17/// A guest opens a section through the host, its edit lands in the host's file, and the
162/// host's own edit reaches the guest.18/// host's own edit reaches the guest.
crates/notebook/tests/support/live.rs created+157
...@@ -0,0 +1,157 @@
1//! What the Live Share tests share: a relay, a host and guests on a notebook of two sections.
2#![allow(dead_code)]
3
4use notebook::{
5 EditStatus, Replica,
6 live::{
7 Hello,
8 share::{self, Guest, Host, Sharing},
9 },
10 session::{Notebook, Section},
11};
12use onestore::{
13 ExGuid,
14 op::{Edit, Op, PageOp},
15 protected::{Key, rekey},
16};
17use std::{
18 net::TcpListener,
19 path::Path,
20 sync::Arc,
21 thread,
22 time::{Duration, Instant},
23};
24
25#[path = "server.rs"]
26pub mod server;
27
28pub const PASSWORD: &str = "fixture password";
29
30/// A relay on this computer with `config`'s limits: its URL.
31pub fn relay(config: relay::server::Config) -> String {
32 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
33 let url = format!("ws://{}", listener.local_addr().unwrap());
34 thread::spawn(move || relay::server::serve(listener, config));
35 url
36}
37
38/// Another secret than `secret`, as a guess makes one.
39pub fn mistaken(secret: &str) -> String {
40 let first = if secret.starts_with('A') { 'B' } else { 'A' };
41 format!("{first}{}", &secret[1..])
42}
43
44pub fn hello(name: &str) -> Hello {
45 Hello::new(name.into(), None).unwrap()
46}
47
48pub fn until(what: &str, done: impl Fn() -> bool) {
49 let deadline = Instant::now() + Duration::from_secs(30);
50 while !done() {
51 assert!(Instant::now() < deadline, "{what}");
52 thread::sleep(Duration::from_millis(20));
53 }
54}
55
56/// A notebook folder holding `Garden.one`, a page reading "Original text", and a protected
57/// `Sealed.one` reading "Sealed text".
58pub fn notebook(root: &Path) -> std::path::PathBuf {
59 let folder = root.join("Garden");
60 std::fs::create_dir_all(&folder).unwrap();
61 std::fs::write(
62 folder.join("Garden.one"),
63 onestore::create_section("Garden.one", "Original text", "Fixture").unwrap(),
64 )
65 .unwrap();
66 let plain = onestore::create_section("Sealed.one", "Sealed text", "Fixture").unwrap();
67 let key = Key::new(PASSWORD).unwrap();
68 std::fs::write(
69 folder.join("Sealed.one"),
70 rekey(&plain, None, Some(&key)).unwrap(),
71 )
72 .unwrap();
73 folder
74}
75
76pub fn host(folder: &Path, cache: &Path, sharing: &Sharing, url: &str) -> Host {
77 let storage = Notebook::open(folder, cache).unwrap().into_storage();
78 Host::start(
79 storage,
80 hello("Ada"),
81 sharing.clone(),
82 "Garden",
83 None,
84 Some(url),
85 || {},
86 )
87 .unwrap()
88}
89
90/// The host's code once the relay has numbered it.
91pub fn code(host: &Host) -> String {
92 until("the code was never numbered", || {
93 host.code().is_some_and(|code| share::code(&code).is_some())
94 });
95 host.code().unwrap()
96}
97
98/// `name` joins with `code` and opens the notebook in `cache`.
99pub fn guest(name: &str, code: &str, url: &str, cache: &Path) -> (Arc<Guest>, Notebook) {
100 let welcome = share::join(hello(name), code, "", None, Some(url)).unwrap();
101 assert_eq!(
102 (welcome.notebook.as_str(), welcome.host.as_str()),
103 ("Garden", "Ada")
104 );
105 let guest = Guest::start(
106 hello(name),
107 welcome.share,
108 welcome.secret,
109 None,
110 Some(url),
111 || {},
112 )
113 .unwrap();
114 until("the host was never met", || guest.host().is_some());
115 let notebook = Notebook::open_hosted(Arc::clone(&guest), cache).unwrap();
116 (guest, notebook)
117}
118
119pub fn open(notebook: &Notebook, guest: &Arc<Guest>, path: &str, key: Option<&Key>) -> Section {
120 let replica = notebook.replica_path(path).unwrap();
121 std::fs::create_dir_all(replica.parent().unwrap()).unwrap();
122 let replica = Replica::open_or_create(&replica, key, || notebook.read_section(path)).unwrap();
123 Section::resume_hosted(path.into(), replica, Arc::clone(guest), || {}).unwrap()
124}
125
126pub fn replace(section: &Section, image: &[u8], range: std::ops::Range<u32>, with: &str) -> u64 {
127 let (space, text, _) = server::text(image);
128 replaced(section, space, text, range, with)
129}
130
131pub fn replaced(
132 section: &Section,
133 space: ExGuid,
134 text: ExGuid,
135 range: std::ops::Range<u32>,
136 with: &str,
137) -> u64 {
138 let op = PageOp::Text {
139 text,
140 range,
141 with: with.into(),
142 };
143 let edit = Edit {
144 at: 134_000_000_000_000_000,
145 ops: vec![Op::Page { space, op }],
146 };
147 section.replica().apply("Guest", edit).unwrap()
148}
149
150pub fn published(section: &Section, id: u64) {
151 until("the edit was never published", || {
152 matches!(
153 section.status(id).unwrap(),
154 Some(EditStatus::Published { .. })
155 )
156 });
157}
crates/relay/README.md+9
...@@ -68,6 +68,15 @@ Share names another relay....@@ -68,6 +68,15 @@ Share names another relay.
68could otherwise claim any address. Listening on a public address without TLS works but lets68could otherwise claim any address. Listening on a public address without TLS works but lets
69anyone on the path see room tags and nameplates.69anyone on the path see room tags and nameplates.
7070
71## Where WebSockets don't get through
72
73Some networks' proxies refuse WebSockets. A peer then joins with `?poll=1` on the same path,
74and the relay answers `session <id>`; `GET /v1/poll/<id>` then waits up to 25 seconds for what
75is to go to it, and `POST /v1/poll/<id>` brings what it sends, each body a run of WebSocket
76frames as the socket would carry them, so the room sees no difference. A session that asks
77nothing for a minute leaves. Through Caddy or nginx this needs nothing more than the
78WebSocket's own proxying; keep nginx's `proxy_read_timeout` above 25 seconds.
79
71## The site: snowbound.paperclover.net80## The site: snowbound.paperclover.net
7281
73`snowbound-site` serves the hosted web build's folder, and for a path that is a Live Share82`snowbound-site` serves the hosted web build's folder, and for a path that is a Live Share
crates/relay/src/server.rs+221-42
...@@ -65,6 +65,11 @@ impl Default for Config {...@@ -65,6 +65,11 @@ impl Default for Config {
65}65}
6666
67const HANDSHAKE: Duration = Duration::from_secs(10);67const HANDSHAKE: Duration = Duration::from_secs(10);
68/// How long a polled session's `GET` waits for something to bring, how long one of its
69/// connections may idle between requests, and how long a session may ask nothing.
70const WAIT: Duration = Duration::from_secs(25);
71const KEPT: Duration = Duration::from_secs(60);
72const IDLE_POLL: Duration = Duration::from_secs(60);
68const WRITE: Duration = Duration::from_secs(30);73const WRITE: Duration = Duration::from_secs(30);
69const MINUTE: Duration = Duration::from_secs(60);74const MINUTE: Duration = Duration::from_secs(60);
70const HOUR: Duration = Duration::from_secs(3600);75const HOUR: Duration = Duration::from_secs(3600);
...@@ -122,6 +127,17 @@ struct Relay {...@@ -122,6 +127,17 @@ struct Relay {
122struct State {127struct State {
123 rooms: HashMap<String, Room>,128 rooms: HashMap<String, Room>,
124 addresses: HashMap<IpAddr, Address>,129 addresses: HashMap<IpAddr, Address>,
130 /// Peers that reach the relay by requests rather than a WebSocket, by session.
131 polls: HashMap<String, Poll>,
132}
133
134/// A peer in a room by requests: where its messages wait for its next `GET`.
135struct Poll {
136 tag: String,
137 slot: u32,
138 outbox: Arc<Outbox>,
139 /// When it last asked anything; one quiet past `IDLE_POLL` has gone.
140 last: Instant,
125}141}
126142
127struct Room {143struct Room {
...@@ -169,6 +185,8 @@ enum Refusal {...@@ -169,6 +185,8 @@ enum Refusal {
169}185}
170186
171impl Relay {187impl Relay {
188 /// Answers a connection's requests one after another, as a polled session sends them,
189 /// until one takes the connection over as a WebSocket or it ends.
172 fn connection(&self, stream: TcpStream) {190 fn connection(&self, stream: TcpStream) {
173 let _ = stream.set_nodelay(true);191 let _ = stream.set_nodelay(true);
174 let _ = stream.set_read_timeout(Some(HANDSHAKE));192 let _ = stream.set_read_timeout(Some(HANDSHAKE));
...@@ -177,16 +195,27 @@ impl Relay {...@@ -177,16 +195,27 @@ impl Relay {
177 return;195 return;
178 };196 };
179 let mut reader = BufReader::new(reading);197 let mut reader = BufReader::new(reading);
180 let Ok(head) = ws::head(&mut reader) else {198 while let Ok(head) = ws::head(&mut reader) {
181 return;199 if !self.request(&stream, &mut reader, &head) {
182 };200 return;
183 let address = self.address(&stream, &head);201 }
202 // Between a session's requests, a connection may idle as long as one waits.
203 let _ = stream.set_read_timeout(Some(KEPT));
204 }
205 }
206
207 /// Answers the request `head`: whether the connection serves another.
208 fn request(&self, stream: &TcpStream, reader: &mut BufReader<TcpStream>, head: &str) -> bool {
209 let address = self.address(stream, head);
184 let target = head.split(' ').nth(1).unwrap_or_default();210 let target = head.split(' ').nth(1).unwrap_or_default();
185 let (path, query) = target.split_once('?').unwrap_or((target, ""));211 let (path, query) = target.split_once('?').unwrap_or((target, ""));
212 if let Some(id) = path.strip_prefix("/v1/poll/") {
213 return self.polled(stream, reader, head, id);
214 }
186 let ask = match path {215 let ask = match path {
187 "/health" => {216 "/health" => {
188 respond(&stream, "200 OK", "application/json", "", &self.health());217 respond(stream, "200 OK", "application/json", "", &self.health());
189 return;218 return false;
190 }219 }
191 "/v1/claim" => Ask::Claim(220 "/v1/claim" => Ask::Claim(
192 query221 query
...@@ -197,53 +226,59 @@ impl Relay {...@@ -197,53 +226,59 @@ impl Relay {
197 _ => match path.strip_prefix("/v1/room/").filter(|tag| valid(tag)) {226 _ => match path.strip_prefix("/v1/room/").filter(|tag| valid(tag)) {
198 Some(tag) => Ask::Room(tag.into()),227 Some(tag) => Ask::Room(tag.into()),
199 None => {228 None => {
200 respond(&stream, "404 Not Found", "text/plain", "", "No such page\n");229 respond(stream, "404 Not Found", "text/plain", "", "No such page\n");
201 return;230 return false;
202 }231 }
203 },232 },
204 };233 };
205 let upgrade = ws::header(&head, "Upgrade")234 if query.split('&').any(|pair| pair == "poll=1") {
235 let outbox = Arc::new(Outbox::mailbox(self.config.queue));
236 return match self.join(address, ask, &outbox, Instant::now()) {
237 Ok((tag, slot)) => {
238 let mut id = [0; 16];
239 if getrandom::fill(&mut id).is_err() {
240 return false;
241 }
242 let id: String = id.iter().map(|byte| format!("{byte:02x}")).collect();
243 let poll = Poll {
244 tag,
245 slot,
246 outbox,
247 last: Instant::now(),
248 };
249 self.state.lock().unwrap().polls.insert(id.clone(), poll);
250 reply(stream, "200 OK", format!("session {id}").as_bytes())
251 }
252 Err(refusal) => {
253 refuse(stream, refusal);
254 false
255 }
256 };
257 }
258 let upgrade = ws::header(head, "Upgrade")
206 .is_some_and(|value| value.eq_ignore_ascii_case("websocket"));259 .is_some_and(|value| value.eq_ignore_ascii_case("websocket"));
207 let (true, Some(key)) = (260 let (true, Some(key)) = (
208 upgrade && head.starts_with("GET "),261 upgrade && head.starts_with("GET "),
209 ws::header(&head, "Sec-WebSocket-Key"),262 ws::header(head, "Sec-WebSocket-Key"),
210 ) else {263 ) else {
211 respond(264 respond(
212 &stream,265 stream,
213 "400 Bad Request",266 "400 Bad Request",
214 "text/plain",267 "text/plain",
215 "",268 "",
216 "A WebSocket only\n",269 "A WebSocket, or ?poll=1\n",
217 );270 );
218 return;271 return false;
219 };272 };
220 let Ok(writing) = stream.try_clone() else {273 let Ok(writing) = stream.try_clone() else {
221 return;274 return false;
222 };275 };
223 let outbox = Arc::new(Outbox::new(writing, self.config.queue));276 let outbox = Arc::new(Outbox::new(writing, self.config.queue));
224 let (tag, slot) = match self.join(address, ask, &outbox, Instant::now()) {277 let (tag, slot) = match self.join(address, ask, &outbox, Instant::now()) {
225 Ok(joined) => joined,278 Ok(joined) => joined,
226 Err(refusal) => {279 Err(refusal) => {
227 let (status, headers, body) = match refusal {280 refuse(stream, refusal);
228 Refusal::NotFound => ("404 Not Found", String::new(), "No such code\n"),281 return false;
229 Refusal::Gone => ("410 Gone", String::new(), "The code has expired\n"),
230 Refusal::Wait(wait) => (
231 "429 Too Many Requests",
232 // Rounded up, so that a client waiting so long finds it over.
233 format!(
234 "Retry-After: {}\r\n",
235 wait.as_secs() + u64::from(wait.subsec_nanos() > 0)
236 ),
237 "Too many tries\n",
238 ),
239 Refusal::Full => (
240 "503 Service Unavailable",
241 "Retry-After: 30\r\n".into(),
242 "The relay is full\n",
243 ),
244 };
245 respond(&stream, status, "text/plain", &headers, body);
246 return;
247 }282 }
248 };283 };
249 let switching = format!(284 let switching = format!(
...@@ -252,7 +287,7 @@ impl Relay {...@@ -252,7 +287,7 @@ impl Relay {
252 ws::accept(key)287 ws::accept(key)
253 );288 );
254 let writer = Arc::clone(&outbox);289 let writer = Arc::clone(&outbox);
255 if (&stream).write_all(switching.as_bytes()).is_ok()290 if (&*stream).write_all(switching.as_bytes()).is_ok()
256 && thread::Builder::new()291 && thread::Builder::new()
257 .stack_size(STACK)292 .stack_size(STACK)
258 .spawn(move || writer.drain())293 .spawn(move || writer.drain())
...@@ -268,6 +303,75 @@ impl Relay {...@@ -268,6 +303,75 @@ impl Relay {
268 }303 }
269 let mut state = self.state.lock().unwrap();304 let mut state = self.state.lock().unwrap();
270 depart(&mut state, &self.config, &tag, slot, Instant::now());305 depart(&mut state, &self.config, &tag, slot, Instant::now());
306 false
307 }
308
309 /// A polled session's request: `GET` waits for what is to go to it, `POST` brings what
310 /// it sends, each a run of WebSocket frames. Whether the connection serves another.
311 fn polled(
312 &self,
313 stream: &TcpStream,
314 reader: &mut BufReader<TcpStream>,
315 head: &str,
316 id: &str,
317 ) -> bool {
318 let length: usize = ws::header(head, "Content-Length")
319 .and_then(|length| length.parse().ok())
320 .unwrap_or(0);
321 if length > self.config.max_message * 4 {
322 return false;
323 }
324 let mut body = vec![0; length];
325 if std::io::Read::read_exact(reader, &mut body).is_err() {
326 return false;
327 }
328 let session = {
329 let mut state = self.state.lock().unwrap();
330 state.polls.get_mut(id).map(|poll| {
331 poll.last = Instant::now();
332 (poll.tag.clone(), poll.slot, Arc::clone(&poll.outbox))
333 })
334 };
335 let Some((tag, slot, outbox)) = session else {
336 return reply(stream, "410 Gone", b"No such session\n");
337 };
338 if head.starts_with("POST ") {
339 let mut frames = ws::Reader::new(&body[..], self.config.max_message, true);
340 while let Ok(message) = frames.read() {
341 if !self.heard(&tag, slot, &outbox, message) {
342 self.end_poll(id);
343 return reply(stream, "410 Gone", b"Closed\n");
344 }
345 }
346 return reply(stream, "200 OK", b"");
347 }
348 let _ = stream.set_write_timeout(Some(WRITE));
349 match outbox.take(WAIT) {
350 Some(bytes) => {
351 if let Some(poll) = self.state.lock().unwrap().polls.get_mut(id) {
352 poll.last = Instant::now();
353 }
354 reply(stream, "200 OK", &bytes)
355 }
356 None => {
357 self.end_poll(id);
358 reply(stream, "410 Gone", b"Closed\n")
359 }
360 }
361 }
362
363 /// Ends the polled session `id`: it leaves its room.
364 fn end_poll(&self, id: &str) {
365 let mut state = self.state.lock().unwrap();
366 if let Some(poll) = state.polls.remove(id) {
367 depart(
368 &mut state,
369 &self.config,
370 &poll.tag,
371 poll.slot,
372 Instant::now(),
373 );
374 }
271 }375 }
272376
273 /// Where a peer counts: the address it connected from, or its proxy says it did.377 /// Where a peer counts: the address it connected from, or its proxy says it did.
...@@ -313,7 +417,9 @@ impl Relay {...@@ -313,7 +417,9 @@ impl Relay {
313 ) -> Result<(String, u32), Refusal> {417 ) -> Result<(String, u32), Refusal> {
314 let config = &self.config;418 let config = &self.config;
315 let mut state = self.state.lock().unwrap();419 let mut state = self.state.lock().unwrap();
316 let State { rooms, addresses } = &mut *state;420 let State {
421 rooms, addresses, ..
422 } = &mut *state;
317 if !addresses.contains_key(&address) && addresses.len() >= ADDRESSES {423 if !addresses.contains_key(&address) && addresses.len() >= ADDRESSES {
318 return Err(Refusal::Full);424 return Err(Refusal::Full);
319 }425 }
...@@ -484,7 +590,9 @@ impl Relay {...@@ -484,7 +590,9 @@ impl Relay {
484 /// Takes a code's owner's word on the peer in a slot waiting to meet it.590 /// Takes a code's owner's word on the peer in a slot waiting to meet it.
485 fn judge(&self, tag: &str, from: u32, verdict: Verdict) {591 fn judge(&self, tag: &str, from: u32, verdict: Verdict) {
486 let mut state = self.state.lock().unwrap();592 let mut state = self.state.lock().unwrap();
487 let State { rooms, addresses } = &mut *state;593 let State {
594 rooms, addresses, ..
595 } = &mut *state;
488 let Some(room) = rooms.get_mut(tag).filter(|room| room.owner == Some(from)) else {596 let Some(room) = rooms.get_mut(tag).filter(|room| room.owner == Some(from)) else {
489 return;597 return;
490 };598 };
...@@ -537,6 +645,15 @@ impl Relay {...@@ -537,6 +645,15 @@ impl Relay {
537 for (tag, slot) in late {645 for (tag, slot) in late {
538 depart(&mut state, config, &tag, slot, now);646 depart(&mut state, config, &tag, slot, now);
539 }647 }
648 let quiet: Vec<String> = (state.polls.iter())
649 .filter(|(_, poll)| now.saturating_duration_since(poll.last) >= IDLE_POLL)
650 .map(|(id, _)| id.clone())
651 .collect();
652 for id in quiet {
653 if let Some(poll) = state.polls.remove(&id) {
654 depart(&mut state, config, &poll.tag, poll.slot, now);
655 }
656 }
540 state.addresses.retain(|_, address| {657 state.addresses.retain(|_, address| {
541 address.refresh(now);658 address.refresh(now);
542 !address.quiet(config, now)659 !address.quiet(config, now)
...@@ -547,7 +664,9 @@ impl Relay {...@@ -547,7 +664,9 @@ impl Relay {
547/// Takes `slot` out of room `tag` and hangs up on it. One still waiting to meet a code's664/// Takes `slot` out of room `tag` and hangs up on it. One still waiting to meet a code's
548/// owner tried a wrong code; the owner leaving excuses those waiting for it.665/// owner tried a wrong code; the owner leaving excuses those waiting for it.
549fn depart(state: &mut State, config: &Config, tag: &str, slot: u32, now: Instant) {666fn depart(state: &mut State, config: &Config, tag: &str, slot: u32, now: Instant) {
550 let State { rooms, addresses } = state;667 let State {
668 rooms, addresses, ..
669 } = state;
551 let Some(room) = rooms.get_mut(tag) else {670 let Some(room) = rooms.get_mut(tag) else {
552 return;671 return;
553 };672 };
...@@ -613,6 +732,38 @@ fn valid(tag: &str) -> bool {...@@ -613,6 +732,38 @@ fn valid(tag: &str) -> bool {
613 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-')732 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-')
614}733}
615734
735/// Answers a polled session's request, keeping the connection: whether that worked.
736fn reply(mut stream: &TcpStream, status: &str, body: &[u8]) -> bool {
737 let head = format!(
738 "HTTP/1.1 {status}\r\nContent-Type: application/octet-stream\r\nContent-Length: {}\r\n\
739 Cache-Control: no-store\r\n\r\n",
740 body.len()
741 );
742 stream.write_all(head.as_bytes()).is_ok() && stream.write_all(body).is_ok()
743}
744
745fn refuse(stream: &TcpStream, refusal: Refusal) {
746 let (status, headers, body) = match refusal {
747 Refusal::NotFound => ("404 Not Found", String::new(), "No such code\n"),
748 Refusal::Gone => ("410 Gone", String::new(), "The code has expired\n"),
749 Refusal::Wait(wait) => (
750 "429 Too Many Requests",
751 // Rounded up, so that a client waiting so long finds it over.
752 format!(
753 "Retry-After: {}\r\n",
754 wait.as_secs() + u64::from(wait.subsec_nanos() > 0)
755 ),
756 "Too many tries\n",
757 ),
758 Refusal::Full => (
759 "503 Service Unavailable",
760 "Retry-After: 30\r\n".into(),
761 "The relay is full\n",
762 ),
763 };
764 respond(stream, status, "text/plain", &headers, body);
765}
766
616fn respond(mut stream: &TcpStream, status: &str, kind: &str, headers: &str, body: &str) {767fn respond(mut stream: &TcpStream, status: &str, kind: &str, headers: &str, body: &str) {
617 let _ = write!(768 let _ = write!(
618 stream,769 stream,
...@@ -729,7 +880,8 @@ impl Bucket {...@@ -729,7 +880,8 @@ impl Bucket {
729880
730/// What waits to be written to one peer, capped in bytes.881/// What waits to be written to one peer, capped in bytes.
731struct Outbox {882struct Outbox {
732 stream: TcpStream,883 /// Where its frames go; none for a polled session's, which its `GET`s take.
884 stream: Option<TcpStream>,
733 queue: Mutex<Queue>,885 queue: Mutex<Queue>,
734 ready: Condvar,886 ready: Condvar,
735 most: usize,887 most: usize,
...@@ -745,13 +897,35 @@ struct Queue {...@@ -745,13 +897,35 @@ struct Queue {
745impl Outbox {897impl Outbox {
746 fn new(stream: TcpStream, most: usize) -> Self {898 fn new(stream: TcpStream, most: usize) -> Self {
747 Self {899 Self {
748 stream,900 stream: Some(stream),
749 queue: Mutex::default(),901 queue: Mutex::default(),
750 ready: Condvar::new(),902 ready: Condvar::new(),
751 most,903 most,
752 }904 }
753 }905 }
754906
907 fn mailbox(most: usize) -> Self {
908 Self {
909 stream: None,
910 queue: Mutex::default(),
911 ready: Condvar::new(),
912 most,
913 }
914 }
915
916 /// Everything queued, waiting up to `wait` for something; none once closed.
917 fn take(&self, wait: Duration) -> Option<Vec<u8>> {
918 let mut queue = self.queue.lock().unwrap();
919 if queue.frames.is_empty() && !queue.closed {
920 queue = self.ready.wait_timeout(queue, wait).unwrap().0;
921 }
922 if queue.closed {
923 return None;
924 }
925 queue.bytes = 0;
926 Some(queue.frames.drain(..).flatten().collect())
927 }
928
755 /// Queues `frame`, or hangs up on a peer that reads too slowly to take it.929 /// Queues `frame`, or hangs up on a peer that reads too slowly to take it.
756 fn push(&self, frame: Vec<u8>) {930 fn push(&self, frame: Vec<u8>) {
757 let mut queue = self.queue.lock().unwrap();931 let mut queue = self.queue.lock().unwrap();
...@@ -776,7 +950,9 @@ impl Outbox {...@@ -776,7 +950,9 @@ impl Outbox {
776 };950 };
777 self.ready.notify_one();951 self.ready.notify_one();
778 drop(queue);952 drop(queue);
779 let _ = self.stream.shutdown(Shutdown::Both);953 if let Some(stream) = &self.stream {
954 let _ = stream.shutdown(Shutdown::Both);
955 }
780 }956 }
781957
782 /// Writes what is queued until closed.958 /// Writes what is queued until closed.
...@@ -795,7 +971,10 @@ impl Outbox {...@@ -795,7 +971,10 @@ impl Outbox {
795 queue = self.ready.wait(queue).unwrap();971 queue = self.ready.wait(queue).unwrap();
796 }972 }
797 };973 };
798 if (&self.stream).write_all(&frame).is_err() {974 let Some(mut stream) = self.stream.as_ref() else {
975 return;
976 };
977 if stream.write_all(&frame).is_err() {
799 self.close();978 self.close();
800 return;979 return;
801 }980 }
crates/snowbound/src/live.rs+1-1
...@@ -236,7 +236,7 @@ pub(crate) fn join(...@@ -236,7 +236,7 @@ pub(crate) fn join(
236 code: &str,236 code: &str,
237 password: &str,237 password: &str,
238) -> Result<live::wire::Welcome, live::share::Refusal> {238) -> Result<live::wire::Welcome, live::share::Refusal> {
239 let me = hello().map_err(|_| live::share::Refusal::Unreachable)?;239 let me = hello().map_err(|_| live::share::Refusal::Unreachable(live::Trouble::Other))?;
240 live::share::join(me, code, password, reach(), relay().as_deref())240 live::share::join(me, code, password, reach(), relay().as_deref())
241}241}
242242
crates/snowbound/src/share.rs+44-6
...@@ -4,7 +4,7 @@...@@ -4,7 +4,7 @@
4use crate::{Library, State};4use crate::{Library, State};
5use accesskit::Role;5use accesskit::Role;
6use notebook::live::{6use notebook::live::{
7 Relayed,7 Relayed, Trouble,
8 share::{self, Refusal, Sharing},8 share::{self, Refusal, Sharing},
9 wire::Welcome,9 wire::Welcome,
10};10};
...@@ -90,13 +90,43 @@ fn refusal(refusal: &Refusal) -> String {...@@ -90,13 +90,43 @@ fn refusal(refusal: &Refusal) -> String {
90 None => "Too many wrong codes from this network. Try again later.".into(),90 None => "Too many wrong codes from this network. Try again later.".into(),
91 },91 },
92 Refusal::Busy => "The Live Share relay is busy. Try again in a minute.".into(),92 Refusal::Busy => "The Live Share relay is busy. Try again in a minute.".into(),
93 Refusal::Unreachable => {93 Refusal::Unreachable(trouble) => unreachable(*trouble),
94 "Can’t reach the Live Share relay. Check your internet connection.".into()
95 }
96 Refusal::TimedOut => "The computer sharing didn’t answer. Try again.".into(),94 Refusal::TimedOut => "The computer sharing didn’t answer. Try again.".into(),
97 }95 }
98}96}
9997
98/// What to tell someone whose computer can't reach the relay, and what to try.
99fn unreachable(trouble: Trouble) -> String {
100 match trouble {
101 Trouble::Dns => "Can’t find the Live Share relay. Check your internet connection, or \
102 try another network."
103 .into(),
104 Trouble::Unreachable => "Can’t connect to the Live Share relay. A firewall may block \
105 it. Try another network."
106 .into(),
107 Trouble::TimedOut => "The Live Share relay didn’t answer. A firewall may block it. Try \
108 another network."
109 .into(),
110 Trouble::ProxyUnreachable => "Can’t reach your proxy server. Check its address in your \
111 network settings."
112 .into(),
113 Trouble::ProxyAuthentication => "Your proxy server asks for a name and password. Set \
114 HTTPS_PROXY to http://name:password@proxy:port, then \
115 try again."
116 .into(),
117 Trouble::ProxyRefused(status) => format!(
118 "Your proxy server won’t connect to the Live Share relay ({status}). Ask whoever \
119 runs your network to allow relay.snowbound.paperclover.net."
120 ),
121 Trouble::Certificate => "Something on your network replaced the relay’s security \
122 certificate. If your network inspects secure connections, add \
123 its certificate to this computer’s trusted certificates."
124 .into(),
125 Trouble::Blocked => "Your network blocks the Live Share relay. Try another network.".into(),
126 Trouble::Other => "Can’t reach the Live Share relay. Try again in a moment.".into(),
127 }
128}
129
100/// A dialog's frame, titled `title` with the BETA badge.130/// A dialog's frame, titled `title` with the BETA badge.
101fn frame(ui: &mut Ui, id: Id, title: &str) {131fn frame(ui: &mut Ui, id: Id, title: &str) {
102 let theme = ui.theme.clone();132 let theme = ui.theme.clone();
...@@ -346,9 +376,17 @@ impl State {...@@ -346,9 +376,17 @@ impl State {
346 true,376 true,
347 );377 );
348 match host.relayed() {378 match host.relayed() {
349 Relayed::Unreachable | Relayed::Refused(..) => status(379 Relayed::Unreachable(trouble) => status(
380 ui,
381 &format!(
382 "{} Until then, only people on this network can join.",
383 unreachable(trouble)
384 ),
385 ),
386 Relayed::Refused(..) => status(
350 ui,387 ui,
351 "Can’t reach the Live Share relay. Only people on this network can join.",388 "The Live Share relay turned this computer away. Only people on this \
389 network can join.",
352 ),390 ),
353 _ => {}391 _ => {}
354 }392 }
crates/snowbound/src/update.rs+16
...@@ -563,7 +563,23 @@ fn download(path: &str, limit: u64) -> Result<Vec<u8>, String> {...@@ -563,7 +563,23 @@ fn download(path: &str, limit: u64) -> Result<Vec<u8>, String> {
563 let roots = (system.iter().chain(bundled))563 let roots = (system.iter().chain(bundled))
564 .map(|der| Certificate::from_der(der).to_owned())564 .map(|der| Certificate::from_der(der).to_owned())
565 .collect();565 .collect();
566 // The proxy Live Share's relay goes through, the system's included; ureq reads only the
567 // environment's.
568 #[cfg(feature = "live")]
569 let proxy = {
570 let host = BASE.trim_start_matches("https://").split('/').next();
571 notebook::live::proxy::for_host(host.unwrap_or_default(), true).and_then(|proxy| {
572 let credentials = (proxy.credentials.as_ref())
573 .map_or_else(String::new, |(name, password)| {
574 format!("{name}:{password}@")
575 });
576 ureq::Proxy::new(&format!("http://{credentials}{proxy}")).ok()
577 })
578 };
579 #[cfg(not(feature = "live"))]
580 let proxy = ureq::Proxy::try_from_env();
566 let agent: ureq::Agent = ureq::Agent::config_builder()581 let agent: ureq::Agent = ureq::Agent::config_builder()
582 .proxy(proxy)
567 .tls_config(583 .tls_config(
568 TlsConfig::builder()584 TlsConfig::builder()
569 .root_certs(RootCerts::Specific(Arc::new(roots)))585 .root_certs(RootCerts::Specific(Arc::new(roots)))