authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-05 21:48:11-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-05 22:21:54-07:00
log0df51c3995b5c6dc655b42eaff33529f17b56d52
tree330ce6dc023c2901d068e6c0bb3560445f2d582d
parent8b4ea55788421f43ea8cdaa7fe335fc8b444c340
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

fix: reconnect Live Share after uncertain HTTP delivery

A lost POST response or truncated GET can leave encrypted stream delivery unknown. Retire that polling session and reconnect instead of replaying or discarding its bytes. fixes #92 Assisted-by: gpt-6.1-sol

2 files changed, 232 insertions(+), 48 deletions(-)

crates/notebook/src/live/transport.rs+27-48
...@@ -383,7 +383,7 @@ fn poll(address: &Address, path: &str) -> Result<Connection, Failure> {...@@ -383,7 +383,7 @@ fn poll(address: &Address, path: &str) -> Result<Connection, Failure> {
383 session,383 session,
384 address,384 address,
385 at,385 at,
386 http: Some(http),386 http,
387 arrived: Vec::new(),387 arrived: Vec::new(),
388 read: 0,388 read: 0,
389 }),389 }),
...@@ -404,30 +404,26 @@ fn post(session: &Session, address: &Address, at: &str) {...@@ -404,30 +404,26 @@ fn post(session: &Session, address: &Address, at: &str) {
404 }404 }
405 std::mem::take(&mut *pending)405 std::mem::take(&mut *pending)
406 };406 };
407 // A connection the relay or the proxy closed meanwhile is opened again, once.407 if http.is_none() {
408 let sent = (0..2).any(|_| {408 http = Http::open(address).ok();
409 if http.is_none() {409 if let Some(opened) = &http {
410 http = Http::open(address).ok();410 let Ok(tcp) = opened.tcp.try_clone() else {
411 if let (Some(opened), Ok(mut open)) = (&http, session.open.lock()) {
412 open.extend(opened.tcp.try_clone());
413 }
414 }
415 let answered = http
416 .as_mut()
417 .map(|connection| connection.ask(address, "POST", at, &batch));
418 match answered {
419 Some(Ok((200 | 204, ..))) => true,
420 Some(Ok((410, ..))) => {
421 session.close();411 session.close();
422 true412 return;
423 }413 };
424 _ => {414 let mut open = session.open.lock().unwrap();
425 http = None;415 if session.closed.load(Ordering::Acquire) {
426 false416 let _ = tcp.shutdown(Shutdown::Both);
417 return;
427 }418 }
419 open.push(tcp);
428 }420 }
429 });421 }
430 if !sent {422 let answered = http
423 .as_mut()
424 .map(|connection| connection.ask(address, "POST", at, &batch));
425 // A lost HTTP answer may follow delivery; replaying its bytes corrupts the stream.
426 if !matches!(answered, Some(Ok((200 | 204, ..)))) {
431 session.close();427 session.close();
432 return;428 return;
433 }429 }
...@@ -456,7 +452,7 @@ struct PollReader {...@@ -456,7 +452,7 @@ struct PollReader {
456 session: Arc<Session>,452 session: Arc<Session>,
457 address: Arc<Address>,453 address: Arc<Address>,
458 at: String,454 at: String,
459 http: Option<Http>,455 http: Http,
460 arrived: Vec<u8>,456 arrived: Vec<u8>,
461 read: usize,457 read: usize,
462}458}
...@@ -467,23 +463,7 @@ impl Read for PollReader {...@@ -467,23 +463,7 @@ impl Read for PollReader {
467 if self.session.closed.load(Ordering::Acquire) {463 if self.session.closed.load(Ordering::Acquire) {
468 return Ok(0);464 return Ok(0);
469 }465 }
470 if self.http.is_none() {466 let answered = self.http.ask(&self.address, "GET", &self.at, &[]);
471 let opened = Http::open(&self.address).map_err(|failure| match failure {
472 Failure::Trouble(_, error) => error,
473 Failure::Refused(status, _) => io::Error::other(format!("{status}")),
474 })?;
475 self.session
476 .open
477 .lock()
478 .unwrap()
479 .extend(opened.tcp.try_clone());
480 self.http = Some(opened);
481 }
482 let answered =
483 self.http
484 .as_mut()
485 .expect("opened above")
486 .ask(&self.address, "GET", &self.at, &[]);
487 match answered {467 match answered {
488 Ok((200, _, body)) => {468 Ok((200, _, body)) => {
489 self.arrived = body;469 self.arrived = body;
...@@ -494,17 +474,12 @@ impl Read for PollReader {...@@ -494,17 +474,12 @@ impl Read for PollReader {
494 return Ok(0);474 return Ok(0);
495 }475 }
496 Ok((status, ..)) => {476 Ok((status, ..)) => {
477 self.session.close();
497 return Err(io::Error::other(format!("The relay answered {status}")));478 return Err(io::Error::other(format!("The relay answered {status}")));
498 }479 }
499 Err(error) => {480 Err(error) => {
500 // Opened again once; a second failure ends the session.481 self.session.close();
501 if self.http.take().is_none() {482 return Err(error);
502 return Err(error);
503 }
504 self.http = Http::open(&self.address).ok();
505 if self.http.is_none() {
506 return Err(error);
507 }
508 }483 }
509 }484 }
510 }485 }
...@@ -606,3 +581,7 @@ impl Write for TlsWriter {...@@ -606,3 +581,7 @@ impl Write for TlsWriter {
606 Ok(())581 Ok(())
607 }582 }
608}583}
584
585#[cfg(test)]
586#[path = "transport_tests.rs"]
587mod tests;
crates/notebook/src/live/transport_tests.rs created+205
...@@ -0,0 +1,205 @@
1use super::*;
2use std::{net::TcpListener, thread::JoinHandle, time::Instant};
3
4fn request(stream: &mut TcpStream) -> io::Result<(String, Vec<u8>)> {
5 let head = ws::head(stream)?;
6 let length = ws::header(&head, "Content-Length")
7 .and_then(|length| length.parse().ok())
8 .unwrap_or(0);
9 let mut body = vec![0; length];
10 stream.read_exact(&mut body)?;
11 Ok((head, body))
12}
13
14fn answer(stream: &mut TcpStream, status: u16, body: &[u8]) -> io::Result<()> {
15 write!(
16 stream,
17 "HTTP/1.1 {status} Test\r\nContent-Length: {}\r\n\r\n",
18 body.len()
19 )?;
20 stream.write_all(body)
21}
22
23fn polling<T: Send + 'static>(
24 serve: impl FnOnce(TcpListener, TcpStream) -> T + Send + 'static,
25) -> (Connection, JoinHandle<T>) {
26 proxy::use_proxy(Some(None));
27 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
28 let address = parse(&format!("ws://{}", listener.local_addr().unwrap())).unwrap();
29 let serving = thread::spawn(move || {
30 let (mut stream, _) = listener.accept().unwrap();
31 stream
32 .set_read_timeout(Some(Duration::from_secs(2)))
33 .unwrap();
34 let (head, _) = request(&mut stream).unwrap();
35 assert!(head.starts_with("GET /v1/room/test?poll=1 "));
36 answer(&mut stream, 200, b"session test").unwrap();
37 serve(listener, stream)
38 });
39 let connection = match poll(&address, "/v1/room/test") {
40 Ok(connection) => connection,
41 Err(_) => panic!("the local polling relay refused the session"),
42 };
43 (connection, serving)
44}
45
46fn accept(listener: &TcpListener) -> Option<TcpStream> {
47 listener.set_nonblocking(true).unwrap();
48 let until = Instant::now() + Duration::from_millis(300);
49 loop {
50 match listener.accept() {
51 Ok((stream, _)) => {
52 stream.set_nonblocking(false).unwrap();
53 stream
54 .set_read_timeout(Some(Duration::from_secs(2)))
55 .unwrap();
56 return Some(stream);
57 }
58 Err(error) if error.kind() == io::ErrorKind::WouldBlock && Instant::now() < until => {
59 thread::sleep(Duration::from_millis(1));
60 }
61 Err(error) if error.kind() == io::ErrorKind::WouldBlock => return None,
62 Err(error) => panic!("could not accept the polling request: {error}"),
63 }
64 }
65}
66
67#[test]
68fn an_accepted_post_with_a_lost_answer_retires_its_stream_without_replay() {
69 let sent = ws::frame(ws::BINARY, b"sealed edit", Some([1, 2, 3, 4]));
70 let expected = sent.clone();
71 let (mut connection, serving) = polling(move |listener, mut reading| {
72 assert!(
73 request(&mut reading)
74 .unwrap()
75 .0
76 .starts_with("GET /v1/poll/test ")
77 );
78 let mut posting = accept(&listener).expect("the edit never reached the relay");
79 let (head, body) = request(&mut posting).unwrap();
80 assert!(head.starts_with("POST /v1/poll/test "));
81 assert_eq!(body, expected);
82 drop(posting);
83 let replay = accept(&listener).map(|mut stream| {
84 let (_, body) = request(&mut stream).unwrap();
85 answer(&mut stream, 200, b"").unwrap();
86 body
87 });
88 reading
89 .set_read_timeout(Some(Duration::from_millis(300)))
90 .unwrap();
91 let closed = reading.read(&mut [0]).is_ok_and(|length| length == 0);
92 let _ = answer(&mut reading, 410, b"");
93 (replay, closed)
94 });
95 connection.writer.write_all(&sent).unwrap();
96 let read = connection.reader.read(&mut [0]);
97 let writing_closed = connection.writer.write(&sent).is_err();
98 (connection.close)();
99 let (replay, reading_closed) = serving.join().unwrap();
100 assert!(
101 replay.is_none(),
102 "the same sealed bytes were replayed after an unknown outcome"
103 );
104 assert!(reading_closed, "the other tracked connection stayed open");
105 assert!(read.is_err() || read.is_ok_and(|length| length == 0));
106 assert!(writing_closed);
107}
108
109#[test]
110fn a_drained_get_with_a_truncated_answer_closes_both_tracked_connections() {
111 let (mut connection, serving) = polling(|listener, mut reading| {
112 assert!(
113 request(&mut reading)
114 .unwrap()
115 .0
116 .starts_with("GET /v1/poll/test ")
117 );
118 let mut posting = accept(&listener).expect("the outgoing connection never opened");
119 assert!(
120 request(&mut posting)
121 .unwrap()
122 .0
123 .starts_with("POST /v1/poll/test ")
124 );
125 reading
126 .write_all(b"HTTP/1.1 200 Test\r\nContent-Length: 5\r\n\r\nxx")
127 .unwrap();
128 drop(reading);
129 posting
130 .set_read_timeout(Some(Duration::from_millis(300)))
131 .unwrap();
132 let closed = posting.read(&mut [0]).is_ok_and(|length| length == 0);
133 let _ = answer(&mut posting, 200, b"");
134 let retry = accept(&listener).map(|mut stream| {
135 assert!(
136 request(&mut stream)
137 .unwrap()
138 .0
139 .starts_with("GET /v1/poll/test ")
140 );
141 answer(&mut stream, 200, b"later").unwrap();
142 });
143 (retry.is_some(), closed)
144 });
145 connection
146 .writer
147 .write_all(&ws::frame(ws::PING, b"", Some([0; 4])))
148 .unwrap();
149 let started = Instant::now();
150 let read = connection.reader.read(&mut [0; 5]);
151 let elapsed = started.elapsed();
152 let writing_closed = read.is_err() && connection.writer.write(b"x").is_err();
153 (connection.close)();
154 let (retried, posting_closed) = serving.join().unwrap();
155 assert!(
156 read.is_err(),
157 "bytes consumed by the failed GET were silently discarded"
158 );
159 assert!(
160 !retried,
161 "a failed GET reused the existing encrypted stream"
162 );
163 assert!(posting_closed, "the other tracked connection stayed open");
164 assert!(writing_closed);
165 assert!(elapsed < Duration::from_secs(1));
166}
167
168#[test]
169fn repeated_get_failures_return_without_a_reconnect_loop() {
170 let (mut connection, serving) = polling(|listener, mut reading| {
171 assert!(
172 request(&mut reading)
173 .unwrap()
174 .0
175 .starts_with("GET /v1/poll/test ")
176 );
177 drop(reading);
178 let mut attempts = 1;
179 while let Some(mut stream) = accept(&listener) {
180 assert!(
181 request(&mut stream)
182 .unwrap()
183 .0
184 .starts_with("GET /v1/poll/test ")
185 );
186 attempts += 1;
187 if attempts == 4 {
188 answer(&mut stream, 410, b"").unwrap();
189 break;
190 }
191 }
192 attempts
193 });
194 let started = Instant::now();
195 let read = connection.reader.read(&mut [0]);
196 let elapsed = started.elapsed();
197 (connection.close)();
198 let attempts = serving.join().unwrap();
199 assert!(
200 read.is_err(),
201 "request failures were retried inside the same session"
202 );
203 assert_eq!(attempts, 1);
204 assert!(elapsed < Duration::from_secs(1));
205}