| author | |
| committer | |
| log | 1f48c20a72a50fe2b69f1cde84c84b025e52973c |
| tree | b2b511c62bcbd3e013cd4ddceae06e1cfc0f1456 |
| parent | fccad14f94e3d22b269957a42ae9c4d3f615edb0 |
| signature | Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU |
Invalidate cached host stamps across presence connectivity transitions and reconcile missed edits when presence returns while preserving the access session. Release callback locks before notifying listeners.
fixes #96
Assisted-by: gpt-6.1-sol3 files changed, 312 insertions(+), 4 deletions(-)
crates/notebook/src/live/share.rs+37-4| ... | ... | @@ -1003,7 +1003,7 @@ struct Inner { |
| 1003 | 1003 | pending: Mutex<HashMap<u64, crate::task::Answer<Reply>>>, |
| 1004 | 1004 | next: AtomicU64, |
| 1005 | 1005 | /// Where the host's reports of changed files go, while a background watches. |
| 1006 | watch: Mutex<Option<Reports>>, | |
| 1006 | watch: Mutex<Option<Arc<Reports>>>, | |
| 1007 | 1007 | /// Why this device can no longer reach the share. |
| 1008 | 1008 | ended: Mutex<Option<Ended>>, |
| 1009 | 1009 | /// Bytes of chunks asked for and not yet given. |
| ... | ... | @@ -1132,7 +1132,7 @@ impl Guest { |
| 1132 | 1132 | if self.host().is_none() { |
| 1133 | 1133 | return Err(self.offline()); |
| 1134 | 1134 | } |
| 1135 | *self.inner.watch.lock().unwrap() = Some(reports); | |
| 1135 | *self.inner.watch.lock().unwrap() = Some(Arc::new(reports)); | |
| 1136 | 1136 | Ok(()) |
| 1137 | 1137 | } |
| 1138 | 1138 | |
| ... | ... | @@ -1610,7 +1610,8 @@ impl Inner { |
| 1610 | 1610 | kind::TOUCHED if serves(from) => { |
| 1611 | 1611 | if let Ok(touched) = minicbor::decode::<Touched>(body) { |
| 1612 | 1612 | self.touched(&touched); |
| 1613 | if let Some(reports) = &*self.watch.lock().unwrap() { | |
| 1613 | let reports = self.watch.lock().unwrap().clone(); | |
| 1614 | if let Some(reports) = reports { | |
| 1614 | 1615 | reports.touched(&touched.paths); |
| 1615 | 1616 | } |
| 1616 | 1617 | } |
| ... | ... | @@ -1639,6 +1640,7 @@ impl Inner { |
| 1639 | 1640 | return; |
| 1640 | 1641 | } |
| 1641 | 1642 | let inner = Arc::downgrade(self); |
| 1643 | let connected = Mutex::new(None); | |
| 1642 | 1644 | let live = Live::start( |
| 1643 | 1645 | self.me.clone(), |
| 1644 | 1646 | &Room::Notebook(secret), |
| ... | ... | @@ -1646,6 +1648,35 @@ impl Inner { |
| 1646 | 1648 | self.relay.as_deref(), |
| 1647 | 1649 | move |event| { |
| 1648 | 1650 | if let Some(inner) = inner.upgrade() { |
| 1651 | if matches!(event, Event::Changed) { | |
| 1652 | let mut connected = connected.lock().unwrap(); | |
| 1653 | let host = inner | |
| 1654 | .host | |
| 1655 | .lock() | |
| 1656 | .unwrap() | |
| 1657 | .as_ref() | |
| 1658 | .map(|(host, _)| host.peer); | |
| 1659 | let peers = inner | |
| 1660 | .presence | |
| 1661 | .lock() | |
| 1662 | .unwrap() | |
| 1663 | .as_ref() | |
| 1664 | .map(|(_, live)| live.peers()) | |
| 1665 | .unwrap_or_default(); | |
| 1666 | let present = | |
| 1667 | host.filter(|host| peers.iter().any(|peer| peer.hello.peer == *host)); | |
| 1668 | if *connected != present { | |
| 1669 | *connected = present; | |
| 1670 | inner.current.lock().unwrap().clear(); | |
| 1671 | drop(connected); | |
| 1672 | let reports = inner.watch.lock().unwrap().clone(); | |
| 1673 | if present.is_some() | |
| 1674 | && let Some(reports) = reports | |
| 1675 | { | |
| 1676 | reports.touched(&[String::new()]); | |
| 1677 | } | |
| 1678 | } | |
| 1679 | } | |
| 1649 | 1680 | if matches!( |
| 1650 | 1681 | event, |
| 1651 | 1682 | Event::Frame { |
| ... | ... | @@ -1663,11 +1694,13 @@ impl Inner { |
| 1663 | 1694 | Ok(live) => { |
| 1664 | 1695 | live.set_presence(self.here.lock().unwrap().clone()); |
| 1665 | 1696 | *presence = Some((secret, live)); |
| 1697 | drop(presence); | |
| 1666 | 1698 | let mut current = self.current.lock().unwrap(); |
| 1667 | 1699 | let paths: Vec<String> = current.keys().cloned().collect(); |
| 1668 | 1700 | current.clear(); |
| 1669 | 1701 | drop(current); |
| 1670 | if let Some(reports) = &*self.watch.lock().unwrap() { | |
| 1702 | let reports = self.watch.lock().unwrap().clone(); | |
| 1703 | if let Some(reports) = reports { | |
| 1671 | 1704 | reports.touched(&paths); |
| 1672 | 1705 | } |
| 1673 | 1706 | } |
crates/notebook/tests/live_share.rs+141| ... | ... | @@ -53,6 +53,147 @@ fn a_guest_edits_the_host_s_notebook() { |
| 53 | 53 | }); |
| 54 | 54 | } |
| 55 | 55 | |
| 56 | #[test] | |
| 57 | fn presence_reconnect_reconciles_missed_changes_without_reconnecting_access() { | |
| 58 | use notebook::{ | |
| 59 | Remote, | |
| 60 | live::{ | |
| 61 | Caret, Presence, Spot, | |
| 62 | share::{Guest, HostedRemote}, | |
| 63 | }, | |
| 64 | session::Background, | |
| 65 | }; | |
| 66 | use onestore::Stamp; | |
| 67 | use std::sync::atomic::{AtomicUsize, Ordering}; | |
| 68 | ||
| 69 | let directory = tempfile::tempdir().unwrap(); | |
| 70 | let folder = notebook(directory.path()); | |
| 71 | let url = relay(Default::default()); | |
| 72 | let sharing = Sharing::new("").unwrap(); | |
| 73 | let host = host(&folder, &directory.path().join("host"), &sharing, &url); | |
| 74 | let welcome = share::join(hello("Grace"), &code(&host), "", None, Some(&url)).unwrap(); | |
| 75 | let relay = PresenceRelay::new(&url, &welcome.room); | |
| 76 | let guest = Guest::start( | |
| 77 | hello("Grace"), | |
| 78 | welcome.share, | |
| 79 | welcome.secret, | |
| 80 | None, | |
| 81 | Some(&relay.url), | |
| 82 | || {}, | |
| 83 | ) | |
| 84 | .unwrap(); | |
| 85 | until("the access host was never met", || guest.host().is_some()); | |
| 86 | let host_id = guest.host().unwrap().peer; | |
| 87 | until("the presence host was never met", || { | |
| 88 | guest.peers().iter().any(|peer| peer.hello.peer == host_id) | |
| 89 | }); | |
| 90 | let mut notebook = | |
| 91 | Notebook::open_hosted(Arc::clone(&guest), directory.path().join("grace")).unwrap(); | |
| 92 | let section = open(&notebook, &guest, "Garden.one", None); | |
| 93 | let background = Background::hosted(Arc::clone(&guest), || {}).unwrap(); | |
| 94 | background.watch(notebook.replicas()); | |
| 95 | background.hold("Garden.one", &section); | |
| 96 | let storage = notebook.into_storage(); | |
| 97 | until("the section never synced", || { | |
| 98 | section.sync_status().unwrap().synced.is_some() | |
| 99 | }); | |
| 100 | let touched = Arc::new(AtomicUsize::new(0)); | |
| 101 | let heard = Arc::clone(&touched); | |
| 102 | let observed = Arc::clone(&guest); | |
| 103 | background.on_touched(Some(Box::new(move |_| { | |
| 104 | observed.set_presence(Presence::default()); | |
| 105 | HostedRemote::new(&observed, "Garden.one").stamp().unwrap(); | |
| 106 | heard.fetch_add(1, Ordering::SeqCst); | |
| 107 | }))); | |
| 108 | let file = folder.join("Garden.one"); | |
| 109 | let image = std::fs::read(&file).unwrap(); | |
| 110 | let (space, text, _) = server::text(&image); | |
| 111 | let stamp = Stamp::of(&image).unwrap(); | |
| 112 | until("the initial TOUCHED never arrived", || { | |
| 113 | host.touched(&["Garden.one".into()]); | |
| 114 | touched.load(Ordering::SeqCst) > 0 | |
| 115 | }); | |
| 116 | let mut remote = HostedRemote::new(&guest, "Garden.one"); | |
| 117 | until("TOUCHED did not cache the current stamp", || { | |
| 118 | let before = relay.state.lock().unwrap().access_bytes; | |
| 119 | assert_eq!( | |
| 120 | HostedRemote::new(&guest, "Garden.one").stamp().unwrap(), | |
| 121 | stamp | |
| 122 | ); | |
| 123 | relay.state.lock().unwrap().access_bytes == before | |
| 124 | }); | |
| 125 | let before = relay.state.lock().unwrap().access_bytes; | |
| 126 | for at in 1..=8 { | |
| 127 | let spot = Spot { | |
| 128 | text: text.into(), | |
| 129 | offset: at, | |
| 130 | }; | |
| 131 | let presence = Presence { | |
| 132 | page: Some(space.into()), | |
| 133 | caret: Some(Caret { | |
| 134 | anchor: spot, | |
| 135 | focus: spot, | |
| 136 | }), | |
| 137 | ..Presence::default() | |
| 138 | }; | |
| 139 | host.set_presence(presence.clone()); | |
| 140 | until("ordinary presence never reached the guest", || { | |
| 141 | guest | |
| 142 | .peers() | |
| 143 | .iter() | |
| 144 | .any(|peer| peer.hello.peer == host_id && peer.presence.as_ref() == Some(&presence)) | |
| 145 | }); | |
| 146 | assert_eq!(remote.stamp().unwrap(), stamp); | |
| 147 | } | |
| 148 | assert_eq!( | |
| 149 | relay.state.lock().unwrap().access_bytes, | |
| 150 | before, | |
| 151 | "typing caused remote stamp or image requests" | |
| 152 | ); | |
| 153 | ||
| 154 | relay.disconnect_presence(); | |
| 155 | until("the presence host never left", || { | |
| 156 | !guest.peers().iter().any(|peer| peer.hello.peer == host_id) | |
| 157 | }); | |
| 158 | assert_eq!(guest.host().unwrap().peer, host_id); | |
| 159 | let access_connections = relay.state.lock().unwrap().access_connections; | |
| 160 | assert_eq!(access_connections, 1); | |
| 161 | std::fs::write(folder.join("reachable.bin"), b"still connected").unwrap(); | |
| 162 | assert_eq!( | |
| 163 | storage.read_file("reachable.bin", 1024).unwrap(), | |
| 164 | b"still connected" | |
| 165 | ); | |
| 166 | let missed = server::typed(&image, space, text, 0..8, "Missed"); | |
| 167 | std::fs::write(&file, &missed).unwrap(); | |
| 168 | let before = touched.load(Ordering::SeqCst); | |
| 169 | host.touched(&["Garden.one".into()]); | |
| 170 | std::thread::sleep(std::time::Duration::from_millis(200)); | |
| 171 | assert_eq!( | |
| 172 | touched.load(Ordering::SeqCst), | |
| 173 | before, | |
| 174 | "the isolated guest received TOUCHED" | |
| 175 | ); | |
| 176 | assert!( | |
| 177 | server::page_texts(&section.page(space).unwrap()).contains(&"Original text".to_owned()) | |
| 178 | ); | |
| 179 | ||
| 180 | relay.state.lock().unwrap().blocked = false; | |
| 181 | until("the presence host never returned", || { | |
| 182 | guest.peers().iter().any(|peer| peer.hello.peer == host_id) | |
| 183 | }); | |
| 184 | until( | |
| 185 | "the missed change never reconciled after presence returned", | |
| 186 | || server::page_texts(&section.page(space).unwrap()).contains(&"Missed text".to_owned()), | |
| 187 | ); | |
| 188 | assert_eq!(guest.host().unwrap().peer, host_id); | |
| 189 | assert_eq!( | |
| 190 | relay.state.lock().unwrap().access_connections, | |
| 191 | access_connections | |
| 192 | ); | |
| 193 | assert_eq!(section.replica().snapshot().unwrap(), missed); | |
| 194 | background.stop(); | |
| 195 | } | |
| 196 | ||
| 56 | 197 | /// While the host is away a guest's edits wait in its replica, the notebook opens from its |
| 57 | 198 | /// last listing, and once the host is back the edits reach its file. |
| 58 | 199 | #[test] |
crates/notebook/tests/support/live.rs+134| ... | ... | @@ -156,3 +156,137 @@ pub fn published(section: &Section, id: u64) { |
| 156 | 156 | ) |
| 157 | 157 | }); |
| 158 | 158 | } |
| 159 | ||
| 160 | pub struct PresenceRelay { | |
| 161 | pub url: String, | |
| 162 | pub state: Arc<std::sync::Mutex<RelayState>>, | |
| 163 | thread: Option<thread::JoinHandle<()>>, | |
| 164 | } | |
| 165 | ||
| 166 | #[derive(Default)] | |
| 167 | pub struct RelayState { | |
| 168 | pub blocked: bool, | |
| 169 | pub access_connections: usize, | |
| 170 | pub access_bytes: usize, | |
| 171 | sockets: Vec<(bool, std::net::TcpStream)>, | |
| 172 | stopped: bool, | |
| 173 | } | |
| 174 | ||
| 175 | impl PresenceRelay { | |
| 176 | pub fn new(upstream: &str, secret: &[u8; 16]) -> Self { | |
| 177 | use sha2::{Digest, Sha256}; | |
| 178 | use std::{ | |
| 179 | io::{self, Read, Write}, | |
| 180 | net::{Shutdown, TcpStream}, | |
| 181 | sync::Mutex, | |
| 182 | }; | |
| 183 | ||
| 184 | let listener = TcpListener::bind("127.0.0.1:0").unwrap(); | |
| 185 | let url = format!("ws://{}", listener.local_addr().unwrap()); | |
| 186 | let upstream = upstream.strip_prefix("ws://").unwrap().to_owned(); | |
| 187 | let tag: String = Sha256::digest([&b"Snowbound room "[..], secret].concat())[..8] | |
| 188 | .iter() | |
| 189 | .map(|byte| format!("{byte:02x}")) | |
| 190 | .collect(); | |
| 191 | let path = format!("/v1/room/{tag}"); | |
| 192 | let state = Arc::new(Mutex::new(RelayState::default())); | |
| 193 | let serving = Arc::clone(&state); | |
| 194 | let thread = thread::spawn(move || { | |
| 195 | let mut workers = Vec::new(); | |
| 196 | for client in listener.incoming().flatten() { | |
| 197 | if serving.lock().unwrap().stopped { | |
| 198 | break; | |
| 199 | } | |
| 200 | let state = Arc::clone(&serving); | |
| 201 | let upstream = upstream.clone(); | |
| 202 | let path = path.clone(); | |
| 203 | workers.push(thread::spawn(move || { | |
| 204 | let _ = (|| -> io::Result<()> { | |
| 205 | let mut client = client; | |
| 206 | client.set_read_timeout(Some(Duration::from_secs(5)))?; | |
| 207 | let head = relay::ws::head(&mut client)?; | |
| 208 | let presence = head | |
| 209 | .split(' ') | |
| 210 | .nth(1) | |
| 211 | .map(|target| target.split('?').next().unwrap()) | |
| 212 | == Some(path.as_str()); | |
| 213 | let mut upstream = TcpStream::connect(upstream)?; | |
| 214 | { | |
| 215 | let mut state = state.lock().unwrap(); | |
| 216 | if state.stopped { | |
| 217 | return Ok(()); | |
| 218 | } | |
| 219 | if presence && state.blocked { | |
| 220 | client.write_all(b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\n\r\n")?; | |
| 221 | return Ok(()); | |
| 222 | } | |
| 223 | state.sockets.push((presence, client.try_clone()?)); | |
| 224 | state.sockets.push((presence, upstream.try_clone()?)); | |
| 225 | if !presence { | |
| 226 | state.access_connections += 1; | |
| 227 | } | |
| 228 | } | |
| 229 | upstream.write_all(head.as_bytes())?; | |
| 230 | client.set_read_timeout(None)?; | |
| 231 | let (mut from, mut to) = (client.try_clone()?, upstream.try_clone()?); | |
| 232 | let counted = Arc::clone(&state); | |
| 233 | let requests = thread::spawn(move || { | |
| 234 | let mut bytes = [0; 16 << 10]; | |
| 235 | while let Ok(length) = from.read(&mut bytes) { | |
| 236 | if length == 0 { | |
| 237 | break; | |
| 238 | } | |
| 239 | if !presence { | |
| 240 | counted.lock().unwrap().access_bytes += length; | |
| 241 | } | |
| 242 | if to.write_all(&bytes[..length]).is_err() { | |
| 243 | break; | |
| 244 | } | |
| 245 | } | |
| 246 | let _ = to.shutdown(Shutdown::Both); | |
| 247 | }); | |
| 248 | let _ = io::copy(&mut upstream, &mut client); | |
| 249 | let _ = client.shutdown(Shutdown::Both); | |
| 250 | let _ = requests.join(); | |
| 251 | Ok(()) | |
| 252 | })(); | |
| 253 | })); | |
| 254 | } | |
| 255 | for worker in workers { | |
| 256 | let _ = worker.join(); | |
| 257 | } | |
| 258 | }); | |
| 259 | Self { | |
| 260 | url, | |
| 261 | state, | |
| 262 | thread: Some(thread), | |
| 263 | } | |
| 264 | } | |
| 265 | ||
| 266 | pub fn disconnect_presence(&self) { | |
| 267 | let mut state = self.state.lock().unwrap(); | |
| 268 | state.blocked = true; | |
| 269 | state.sockets.retain(|(presence, socket)| { | |
| 270 | if *presence { | |
| 271 | let _ = socket.shutdown(std::net::Shutdown::Both); | |
| 272 | } | |
| 273 | !*presence | |
| 274 | }); | |
| 275 | } | |
| 276 | } | |
| 277 | ||
| 278 | impl Drop for PresenceRelay { | |
| 279 | fn drop(&mut self) { | |
| 280 | { | |
| 281 | let mut state = self.state.lock().unwrap(); | |
| 282 | state.stopped = true; | |
| 283 | for (_, socket) in &state.sockets { | |
| 284 | let _ = socket.shutdown(std::net::Shutdown::Both); | |
| 285 | } | |
| 286 | } | |
| 287 | let _ = std::net::TcpStream::connect(self.url.strip_prefix("ws://").unwrap()); | |
| 288 | if let Some(thread) = self.thread.take() { | |
| 289 | let _ = thread.join(); | |
| 290 | } | |
| 291 | } | |
| 292 | } |