| 1 | use super::*; |
| 2 | use crate::task::{self, JoinHandle}; |
| 3 | use std::{ |
| 4 | collections::hash_map::RandomState, |
| 5 | hash::BuildHasher, |
| 6 | sync::{ |
| 7 | Arc, Mutex, |
| 8 | atomic::{AtomicBool, AtomicU64, Ordering}, |
| 9 | }, |
| 10 | }; |
| 11 | use web_time::Instant; |
| 12 | |
| 13 | pub(super) struct Signal { |
| 14 | pub(crate) stopped: AtomicBool, |
| 15 | pub(crate) offline: AtomicBool, |
| 16 | /// FILETIME of the last step or poll that reached the remote; 0 before one has. |
| 17 | synced: AtomicU64, |
| 18 | /// Asked for by `wake`: steps run while offline until one leaves nothing to do. |
| 19 | pub(crate) requested: AtomicBool, |
| 20 | /// A watch on the file's folder wakes the worker when the file changes, so an idle worker |
| 21 | /// waits for that instead of checking the file on its own (`Background::hold`). |
| 22 | pub(crate) watched: AtomicBool, |
| 23 | pause: Mutex<Pause>, |
| 24 | sender: task::Sender<()>, |
| 25 | } |
| 26 | |
| 27 | /// How local edits wait before they publish (`SyncWorker::set_pause`). |
| 28 | #[derive(Default)] |
| 29 | struct Pause { |
| 30 | length: Duration, |
| 31 | /// When the first and the last edit waiting arrived. |
| 32 | waiting: Option<(Instant, Instant)>, |
| 33 | } |
| 34 | |
| 35 | /// The longest local edits wait for a pause in typing. |
| 36 | const LONGEST: Duration = Duration::from_secs(30); |
| 37 | |
| 38 | impl Signal { |
| 39 | pub(crate) fn new() -> (Arc<Self>, task::Receiver<()>) { |
| 40 | let (sender, receiver) = task::channel(); |
| 41 | let signal = Self { |
| 42 | stopped: AtomicBool::new(false), |
| 43 | offline: AtomicBool::new(false), |
| 44 | synced: AtomicU64::new(0), |
| 45 | requested: AtomicBool::new(false), |
| 46 | watched: AtomicBool::new(false), |
| 47 | pause: Mutex::default(), |
| 48 | sender, |
| 49 | }; |
| 50 | (Arc::new(signal), receiver) |
| 51 | } |
| 52 | |
| 53 | pub(super) fn wake(&self) { |
| 54 | // One retained notification covers edits that arrive during network I/O. |
| 55 | let _ = self.sender.try_send(()); |
| 56 | } |
| 57 | |
| 58 | /// A burst of local edits is durable. |
| 59 | pub(super) fn edited(&self) { |
| 60 | if let Ok(mut pause) = self.pause.lock() { |
| 61 | let now = Instant::now(); |
| 62 | pause.waiting = Some((pause.waiting.map_or(now, |(first, _)| first), now)); |
| 63 | } |
| 64 | self.wake(); |
| 65 | } |
| 66 | |
| 67 | /// How much longer waiting edits wait for a pause in typing, if they do. |
| 68 | fn paused(&self) -> Option<Duration> { |
| 69 | let pause = self.pause.lock().ok()?; |
| 70 | let (first, last) = pause.waiting?; |
| 71 | let now = Instant::now(); |
| 72 | Some( |
| 73 | (last + pause.length) |
| 74 | .min(first + LONGEST) |
| 75 | .checked_duration_since(now)?, |
| 76 | ) |
| 77 | .filter(|wait| !wait.is_zero()) |
| 78 | } |
| 79 | } |
| 80 | |
| 81 | /// Owns automatic reconciliation. Dropping requests cancellation without blocking. |
| 82 | /// The in-flight sync step retains cache ownership until it finishes. |
| 83 | pub struct SyncWorker { |
| 84 | signal: Arc<Signal>, |
| 85 | thread: Option<JoinHandle<Result<()>>>, |
| 86 | } |
| 87 | |
| 88 | impl SyncWorker { |
| 89 | /// Requests a retry, for example after a network reachability change. Working |
| 90 | /// offline, it synchronizes once, as OneNote's Sync Now does. |
| 91 | /// A pending contention backoff finishes before processing the notification. |
| 92 | pub fn wake(&self) { |
| 93 | self.signal.requested.store(true, Ordering::Release); |
| 94 | self.signal.wake(); |
| 95 | } |
| 96 | |
| 97 | /// When a step or poll last reached the remote. |
| 98 | pub fn synced(&self) -> Option<u64> { |
| 99 | Some(self.signal.synced.load(Ordering::Acquire)).filter(|time| *time != 0) |
| 100 | } |
| 101 | |
| 102 | /// Local edits publish once `pause` passes without another, or once they have waited |
| 103 | /// half a minute, rather than each burst at once: a cloud drive uploads every publication |
| 104 | /// and turns each concurrent one into a conflict version. `wake` publishes them at once. |
| 105 | pub fn set_pause(&self, pause: Duration) { |
| 106 | if let Ok(mut current) = self.signal.pause.lock() { |
| 107 | current.length = pause; |
| 108 | } |
| 109 | self.signal.wake(); |
| 110 | } |
| 111 | |
| 112 | /// Working offline, the worker neither connects nor steps: local edits stay queued |
| 113 | /// until `wake`, or until working online again. |
| 114 | pub fn set_offline(&self, offline: bool) { |
| 115 | self.signal.offline.store(offline, Ordering::Release); |
| 116 | self.signal.wake(); |
| 117 | } |
| 118 | |
| 119 | /// Cancels future steps; native threads finish the current step and callback before returning. |
| 120 | /// A stopped worker leaves pending edits and uncertain attempts in the cache. |
| 121 | /// Call outside the worker's own callback, which cannot join its calling thread. |
| 122 | pub fn stop(mut self) -> Result<()> { |
| 123 | self.signal.stopped.store(true, Ordering::Release); |
| 124 | self.signal.wake(); |
| 125 | let thread = self.thread.take().expect("Worker owns its thread"); |
| 126 | #[cfg(target_arch = "wasm32")] |
| 127 | return thread.finished().unwrap_or(Ok(())); |
| 128 | #[cfg(not(target_arch = "wasm32"))] |
| 129 | thread |
| 130 | .join() |
| 131 | .map_err(|_| io::Error::other("Synchronization worker panicked"))? |
| 132 | } |
| 133 | } |
| 134 | |
| 135 | impl Drop for SyncWorker { |
| 136 | fn drop(&mut self) { |
| 137 | self.signal.stopped.store(true, Ordering::Release); |
| 138 | self.signal.wake(); |
| 139 | } |
| 140 | } |
| 141 | |
| 142 | impl Replica { |
| 143 | /// Starts one worker, reconnecting through `connect` after transport failures. |
| 144 | /// Local edits wake it; `interval` controls idle polling, unless a watch wakes it |
| 145 | /// (`Background::hold`), and transport retries. While |
| 146 | /// nothing is queued, or the queue waits on a remote that has not changed since, a |
| 147 | /// remote whose `stamp` holds is not read again. |
| 148 | /// Contended operations returning `NotCommitted` also back off by up to one second. |
| 149 | /// `observe` runs on the worker after each attempt, including connection errors. |
| 150 | /// Cache/document errors stop the worker; inspect them through `observe` or `stop`. |
| 151 | /// Remote calls and callbacks must be bounded for `stop` to have bounded latency. |
| 152 | pub fn start_sync<R, F, O>( |
| 153 | self: &Arc<Self>, |
| 154 | interval: Duration, |
| 155 | mut connect: F, |
| 156 | mut observe: O, |
| 157 | ) -> io::Result<SyncWorker> |
| 158 | where |
| 159 | R: Remote + 'static, |
| 160 | F: FnMut() -> io::Result<R> + Send + 'static, |
| 161 | O: FnMut(&Result<Synced>) + Send + 'static, |
| 162 | { |
| 163 | if interval.is_zero() || Instant::now().checked_add(interval).is_none() { |
| 164 | return Err(io::Error::new( |
| 165 | io::ErrorKind::InvalidInput, |
| 166 | "Synchronization interval must be positive and representable", |
| 167 | )); |
| 168 | } |
| 169 | let mut owner = self |
| 170 | .section |
| 171 | .worker |
| 172 | .lock() |
| 173 | .map_err(|_| io::Error::other("Synchronization worker registration panicked"))?; |
| 174 | if owner.upgrade().is_some() { |
| 175 | return Err(io::ErrorKind::WouldBlock.into()); |
| 176 | } |
| 177 | let (signal, receiver) = Signal::new(); |
| 178 | let replica = Arc::clone(self); |
| 179 | let worker_signal = Arc::clone(&signal); |
| 180 | let thread = task::spawn("onestore-sync", move || async move { |
| 181 | let mut remote: Option<R> = None; |
| 182 | let jitter = RandomState::new(); |
| 183 | let mut contention = 0_u32; |
| 184 | // The first step, and the first after a failure, always runs, so `observe` |
| 185 | // hears that the remote is reachable and what state the queue is in. |
| 186 | let mut reported = false; |
| 187 | // With nothing to do: a watched worker waits for the watch, any other looks |
| 188 | // again after `interval`. |
| 189 | let rest = || { |
| 190 | task::wait( |
| 191 | &receiver, |
| 192 | (!worker_signal.watched.load(Ordering::Acquire)).then_some(interval), |
| 193 | ) |
| 194 | }; |
| 195 | while !worker_signal.stopped.load(Ordering::Acquire) { |
| 196 | if worker_signal.offline.load(Ordering::Acquire) |
| 197 | && !worker_signal.requested.load(Ordering::Acquire) |
| 198 | { |
| 199 | remote = None; |
| 200 | reported = false; |
| 201 | task::wait(&receiver, None).await; |
| 202 | continue; |
| 203 | } |
| 204 | if !worker_signal.requested.load(Ordering::Acquire) |
| 205 | && let Some(wait) = worker_signal.paused() |
| 206 | { |
| 207 | task::wait(&receiver, Some(wait)).await; |
| 208 | continue; |
| 209 | } |
| 210 | if let Ok(mut pause) = worker_signal.pause.lock() { |
| 211 | // The step publishes what waits; edits from now on wait anew. |
| 212 | pause.waiting = None; |
| 213 | } |
| 214 | let result = match remote.as_mut() { |
| 215 | Some(remote) => { |
| 216 | if reported && replica.settled(remote).await.unwrap_or(false) { |
| 217 | worker_signal.requested.store(false, Ordering::Release); |
| 218 | worker_signal.synced.store(crate::now(), Ordering::Release); |
| 219 | rest().await; |
| 220 | continue; |
| 221 | } |
| 222 | let result = replica.sync_once_async(remote).await; |
| 223 | reported = result.is_ok(); |
| 224 | result |
| 225 | } |
| 226 | None => match connect() { |
| 227 | Ok(connected) => { |
| 228 | remote = Some(connected); |
| 229 | continue; |
| 230 | } |
| 231 | Err(error) => Err(Error::RemoteIo(error)), |
| 232 | }, |
| 233 | }; |
| 234 | if result.is_ok() { |
| 235 | worker_signal.synced.store(crate::now(), Ordering::Release); |
| 236 | } |
| 237 | observe(&result); |
| 238 | let idle = matches!(result, Ok(Synced { edit: None, .. })); |
| 239 | match result { |
| 240 | Ok(Synced { |
| 241 | edit: Some((_, EditStatus::Published { .. })), |
| 242 | .. |
| 243 | }) => { |
| 244 | contention = 0; |
| 245 | continue; |
| 246 | } |
| 247 | Ok(Synced { edit: None, .. }) => contention = 0, |
| 248 | Ok(_) => {} |
| 249 | Err(Error::Remote(onestore::CommitError { |
| 250 | state: onestore::CommitState::NotCommitted, |
| 251 | ref error, |
| 252 | })) if matches!( |
| 253 | error.kind(), |
| 254 | io::ErrorKind::WouldBlock | io::ErrorKind::ResourceBusy |
| 255 | ) => |
| 256 | { |
| 257 | contention = contention.saturating_add(1); |
| 258 | let ceiling = (50_u64 << contention.min(5)).min(1000); |
| 259 | let until = Instant::now() |
| 260 | + Duration::from_millis(jitter.hash_one(contention) % ceiling); |
| 261 | // Local wakes must not keep competing writers in the same retry phase. |
| 262 | while !worker_signal.stopped.load(Ordering::Acquire) { |
| 263 | let Some(remaining) = until.checked_duration_since(Instant::now()) |
| 264 | else { |
| 265 | break; |
| 266 | }; |
| 267 | task::wait(&receiver, Some(remaining)).await; |
| 268 | } |
| 269 | continue; |
| 270 | } |
| 271 | Err(Error::RemoteIo(ref error)) |
| 272 | if matches!( |
| 273 | error.kind(), |
| 274 | io::ErrorKind::WouldBlock | io::ErrorKind::ResourceBusy |
| 275 | ) => {} |
| 276 | Err(Error::RemoteIo(_) | Error::Remote(_)) => remote = None, |
| 277 | Err(Error::Io(ref error)) if error.kind() == io::ErrorKind::WouldBlock => {} |
| 278 | Err(error) => return Err(error), |
| 279 | } |
| 280 | worker_signal.requested.store(false, Ordering::Release); |
| 281 | if idle { |
| 282 | rest().await; |
| 283 | } else if !worker_signal.stopped.load(Ordering::Acquire) { |
| 284 | task::wait(&receiver, Some(interval)).await; |
| 285 | } |
| 286 | } |
| 287 | Ok(()) |
| 288 | })?; |
| 289 | *owner = Arc::downgrade(&signal); |
| 290 | Ok(SyncWorker { |
| 291 | signal, |
| 292 | thread: Some(thread), |
| 293 | }) |
| 294 | } |
| 295 | } |