1use super::*;
2use crate::task::{self, JoinHandle};
3use std::{
4 collections::hash_map::RandomState,
5 hash::BuildHasher,
6 sync::{
7 Arc, Mutex,
8 atomic::{AtomicBool, AtomicU64, Ordering},
9 },
10};
11use web_time::Instant;
12
13pub(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)]
29struct 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.
36const LONGEST: Duration = Duration::from_secs(30);
37
38impl 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.
83pub struct SyncWorker {
84 signal: Arc<Signal>,
85 thread: Option<JoinHandle<Result<()>>>,
86}
87
88impl 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
135impl Drop for SyncWorker {
136 fn drop(&mut self) {
137 self.signal.stopped.store(true, Ordering::Release);
138 self.signal.wake();
139 }
140}
141
142impl 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}