web browser (Servo)
git clone https://git.lucas.co/cce-browser.git
src/raindrop/sync.rs (18.3K)
1 //! The sync worker — phase 3: when a pass runs, and what it says.
2 //!
3 //! One thread, `cce-raindrop`, for the life of the browser once the setting
4 //! has been on. It polls the in-memory bookmarks every [`POLL`] instead of
5 //! being told about edits: a bookmark changes from the star, Ctrl+D, the
6 //! bookmarks menu and the `cce://bookmarks` page, and comparing the rows
7 //! (no I/O, a lock held for a copy) catches all of them without a hook in
8 //! each. A change syncs once the rows have held still for one poll, so a
9 //! burst of edits is one pass; Raindrop's own changes arrive with the pass
10 //! every [`FULL`]; the page's "sync now" asks for one at once.
11 //!
12 //! A pass is [`run_pass`]: fetch, plan, apply Raindrop's half, apply the local
13 //! half through `Bookmarks::edit_rows` (one locked edit — no star can land in
14 //! the middle of it), and save the base last. Its outcome becomes the status
15 //! line on `cce://bookmarks`, and what it changed is appended to the sync log
16 //! (`raindrop-sync.log`, beside the base) — one line per bookmark, so "which
17 //! ones?" has an answer after the status line and the process are gone.
18
19 use std::path::Path;
20 use std::sync::atomic::{AtomicBool, Ordering};
21 use std::sync::Arc;
22 use std::time::{Duration, Instant};
23
24 use super::{api, apply_local, load_base, plan, save_base, settle_base, Local};
25 use crate::pages::{Bookmarks, SyncNote, SYNC_FORCE};
26
27 /// How often the rows are compared, and the stillness a change waits for.
28 const POLL: Duration = Duration::from_secs(2);
29 /// A full pass this often while nothing local changes: Raindrop's side.
30 const FULL: Duration = Duration::from_secs(600);
31
32 /// The sync log beside a base: `raindrop-sync.tsv` → `raindrop-sync.log`.
33 pub fn log_path(base_path: &Path) -> std::path::PathBuf {
34 base_path.with_extension("log")
35 }
36
37 /// Past this size the log keeps only its newer half — months of passes, kept
38 /// to a size nobody has to think about.
39 const LOG_MAX: usize = 512 * 1024;
40
41 /// Append `events` (`(what, link, title)`) to the log, one line each:
42 /// `time \t what \t link \t title`. Best effort — a log that cannot be
43 /// written must not fail a pass that already happened — but said once.
44 pub fn append_log(path: &Path, events: &[(String, String, String)]) {
45 if events.is_empty() {
46 return;
47 }
48 let now = api::iso8601(super::unix_now());
49 let flat = |s: &str| s.replace(['\t', '\n', '\r'], " ");
50 let mut text = std::fs::read_to_string(path).unwrap_or_default();
51 for (what, link, title) in events {
52 text.push_str(&format!("{now}\t{}\t{}\t{}\n", flat(what), flat(link), flat(title)));
53 }
54 if text.len() > LOG_MAX {
55 let cut = text.len() - LOG_MAX / 2;
56 let from = text[cut..].find('\n').map(|i| cut + i + 1).unwrap_or(cut);
57 text = text[from..].to_string();
58 }
59 if let Some(dir) = path.parent() {
60 let _ = std::fs::create_dir_all(dir);
61 }
62 let tmp = path.with_extension("log.tmp");
63 if let Err(e) = std::fs::write(&tmp, text).and_then(|_| std::fs::rename(&tmp, path)) {
64 log::warn!("raindrop: could not write {}: {e}", path.display());
65 }
66 }
67
68 /// What one pass did.
69 #[derive(Debug, Default, PartialEq)]
70 pub struct Summary {
71 pub added_here: usize,
72 pub removed_here: usize,
73 pub changed_here: usize,
74 pub added_there: usize,
75 pub trashed_there: usize,
76 pub renamed_there: usize,
77 /// Local operations left for the next pass (edited mid-pass).
78 pub deferred: usize,
79 /// Raindrop calls that failed; each is retried next pass.
80 pub errors: Vec<String>,
81 }
82
83 impl Summary {
84 /// The status line: what moved, or "in sync".
85 pub fn text(&self) -> String {
86 let mut parts = Vec::new();
87 let mut say = |n: usize, what: &str| {
88 if n > 0 {
89 parts.push(format!("{n} {what}"));
90 }
91 };
92 say(self.added_here, "added here");
93 say(self.removed_here, "removed here");
94 say(self.changed_here, "updated here");
95 say(self.added_there, "added to Raindrop");
96 say(self.trashed_there, "moved to Raindrop's trash");
97 say(self.renamed_there, "renamed in Raindrop");
98 let mut s = if parts.is_empty() { "in sync".to_string() } else { format!("synced — {}", parts.join(", ")) };
99 if !self.errors.is_empty() {
100 s.push_str(&format!(" ({} failed, retrying next pass)", self.errors.len()));
101 }
102 s
103 }
104 }
105
106 #[derive(Debug, PartialEq)]
107 pub enum Outcome {
108 Synced(Summary),
109 /// The guard refused the pass; nothing was changed on either side.
110 Refused(String),
111 }
112
113 fn locals(rows: Vec<(u64, String, String)>) -> Vec<Local> {
114 rows.into_iter().map(|(ts, url, title)| Local { url, title, ts }).collect()
115 }
116
117 /// One pass against `client`'s Unsorted. `force` runs a plan the guard
118 /// refused — only ever from the page's "sync anyway", after a person has read
119 /// why it was refused.
120 pub fn run_pass(
121 client: &api::Client,
122 bookmarks: &Bookmarks,
123 base_path: &Path,
124 force: bool,
125 ) -> Result<Outcome, String> {
126 let remote = client.fetch(api::UNSORTED).map_err(|e| e.to_string())?;
127 let snapshot = locals(bookmarks.rows());
128 let base = load_base(base_path);
129 let log = log_path(base_path);
130 let plan = match plan(&snapshot, &remote, &base) {
131 Ok(p) => p,
132 Err(refused) if force => {
133 log::warn!("raindrop: running a refused pass on request ({})", refused.reason);
134 append_log(&log, &[("forced".into(), String::new(), refused.reason.clone())]);
135 refused.plan
136 }
137 Err(refused) => {
138 append_log(&log, &[("refused".into(), String::new(), refused.reason.clone())]);
139 return Ok(Outcome::Refused(refused.reason));
140 }
141 };
142 let applied = api::apply_remote(client, api::UNSORTED, &plan).map_err(|e| e.to_string())?;
143 let touches_local = !(plan.relink_local.is_empty()
144 && plan.rename_local.is_empty()
145 && plan.delete_local.is_empty()
146 && plan.add_local.is_empty());
147 // Only rewrite the file when there is something to write: a pass every
148 // ten minutes that changes nothing must not touch the disk.
149 // The log says what *happened*: the local lines are the difference the
150 // edit actually made (a deferred operation made none), the Raindrop lines
151 // the calls that succeeded.
152 let mut events: Vec<(String, String, String)> = Vec::new();
153 let skipped = if touches_local {
154 bookmarks.edit_rows(|rows| {
155 let before: Vec<(u64, String, String)> = rows.clone();
156 let mut current = locals(std::mem::take(rows));
157 let skipped = apply_local(&mut current, &snapshot, &plan);
158 *rows = current.into_iter().map(|l| (l.ts, l.url, l.title)).collect();
159 for (_, url, title) in rows.iter() {
160 match before.iter().find(|b| b.1 == *url) {
161 None => events.push(("added here".into(), url.clone(), title.clone())),
162 Some(b) if b.2 != *title => {
163 events.push(("renamed here".into(), url.clone(), format!("{} → {title}", b.2)))
164 }
165 Some(_) => {}
166 }
167 }
168 for (_, url, title) in &before {
169 if !rows.iter().any(|r| r.1 == *url) {
170 events.push(("removed here".into(), url.clone(), title.clone()));
171 }
172 }
173 skipped
174 })
175 } else {
176 super::Skipped::default()
177 };
178 let link_of = |id: &super::RaindropId| {
179 remote.iter().find(|r| r.id == *id).map(|r| (r.link.clone(), r.title.clone())).unwrap_or_default()
180 };
181 for (url, _) in &applied.created {
182 let title = plan.create_remote.iter().find(|l| l.url == *url).map(|l| l.title.clone()).unwrap_or_default();
183 events.push(("added to Raindrop".into(), url.clone(), title));
184 }
185 for id in plan.trash_remote.iter().filter(|id| !applied.failed_trash.contains(id)) {
186 let (link, title) = link_of(id);
187 events.push(("moved to Raindrop's trash".into(), link, title));
188 }
189 for (id, new) in plan.rename_remote.iter().filter(|(id, _)| !applied.failed_renames.contains(id)) {
190 let (link, old) = link_of(id);
191 events.push(("renamed in Raindrop".into(), link, format!("{old} → {new}")));
192 }
193 for e in &applied.errors {
194 events.push(("failed".into(), String::new(), e.clone()));
195 }
196 append_log(&log, &events);
197 // Last: a pass that dies before this point leaves the old base, and the
198 // next pass redoes the work rather than misreading it.
199 let new_base = settle_base(&plan, &base, &applied, &skipped);
200 if new_base != base {
201 save_base(base_path, &new_base).map_err(|e| format!("could not save the sync state: {e}"))?;
202 }
203 Ok(Outcome::Synced(Summary {
204 added_here: plan.add_local.len(),
205 removed_here: plan.delete_local.len(),
206 changed_here: plan.rename_local.len() + plan.relink_local.len(),
207 added_there: applied.created.len(),
208 trashed_there: plan.trash_remote.len() - applied.failed_trash.len(),
209 renamed_there: plan.rename_remote.len() - applied.failed_renames.len(),
210 deferred: skipped.count,
211 errors: applied.errors,
212 }))
213 }
214
215 /// A code for the "sync anyway" link. Not cryptographic, and it needs not be:
216 /// it only has to be something a web page cannot know, and the page that
217 /// shows it is the only place it is written.
218 fn force_code() -> u64 {
219 use std::hash::{BuildHasher, Hasher};
220 let mut h = std::collections::hash_map::RandomState::new().build_hasher();
221 h.write_u64(super::unix_now());
222 h.finish()
223 }
224
225 /// One pass, start to status line: the token, the client, `run_pass`.
226 fn pass(bookmarks: &Bookmarks, force: bool) {
227 let note = |text: String, force_code: Option<u64>| {
228 bookmarks.set_sync_note(Some(SyncNote { text, at: super::unix_now(), force_code }));
229 };
230 let token = match crate::accounts::raindrop_token() {
231 Ok(t) => t,
232 Err(e) => {
233 log::warn!("raindrop: {e}");
234 note(e, None);
235 return;
236 }
237 };
238 let client = api::Client::new(token);
239 match run_pass(&client, bookmarks, &super::state_path(), force) {
240 Ok(Outcome::Synced(s)) => {
241 for e in &s.errors {
242 log::warn!("raindrop: {e}");
243 }
244 log::info!("raindrop: {}", s.text());
245 note(s.text(), None);
246 }
247 Ok(Outcome::Refused(reason)) => {
248 log::warn!("raindrop: pass refused: {reason}");
249 // The same code for as long as passes keep being refused: a page
250 // already showing "sync anyway" must stay able to use it when a
251 // later pass — the periodic one, a reload — refuses again.
252 let code = bookmarks.sync_force_code().unwrap_or_else(force_code);
253 note(format!("not synced: {reason}"), Some(code));
254 }
255 Err(e) => {
256 log::warn!("raindrop: {e}");
257 append_log(&log_path(&super::state_path()), &[("not synced".into(), String::new(), e.clone())]);
258 note(format!("not synced: {e}"), None);
259 }
260 }
261 }
262
263 /// Start the worker. `enabled` is the setting, live: off, the worker idles
264 /// and the page shows no status line; on again, it syncs at once.
265 pub fn spawn(bookmarks: Arc<Bookmarks>, enabled: Arc<AtomicBool>) {
266 std::thread::Builder::new()
267 .name("cce-raindrop".to_string())
268 .spawn(move || {
269 let mut last_rows = None;
270 let mut changed = false;
271 let mut next_full = Instant::now();
272 let mut was_on = false;
273 loop {
274 std::thread::sleep(POLL);
275 let on = enabled.load(Ordering::SeqCst);
276 if !on {
277 if was_on {
278 bookmarks.set_sync_note(None);
279 }
280 was_on = false;
281 continue;
282 }
283 if !was_on {
284 // Just turned on (or launched on): sync now.
285 was_on = true;
286 next_full = Instant::now();
287 }
288 let request = bookmarks.take_sync_request();
289 let rows = bookmarks.rows();
290 if last_rows.as_ref() != Some(&rows) {
291 // Changed since the last look: wait one poll for it to
292 // hold still, so a burst of edits is one pass.
293 changed = last_rows.is_some();
294 last_rows = Some(rows);
295 if request == 0 {
296 continue;
297 }
298 }
299 if request != 0 || changed || Instant::now() >= next_full {
300 pass(&bookmarks, request == SYNC_FORCE);
301 changed = false;
302 next_full = Instant::now() + FULL;
303 // The pass's own edits are not a local change to sync.
304 last_rows = Some(bookmarks.rows());
305 }
306 }
307 })
308 .expect("spawn the Raindrop sync worker");
309 }
310
311 #[cfg(test)]
312 mod tests {
313 use super::*;
314 use crate::accounts::Secret;
315
316 #[test]
317 fn the_log_keeps_its_newer_half() {
318 let dir = std::env::temp_dir().join(format!("cce-raindrop-log-{}", std::process::id()));
319 let path = dir.join("raindrop-sync.log");
320 let long = "x".repeat(1000);
321 for i in 0..700 {
322 append_log(&path, &[("added here".into(), format!("https://{i}.test/"), long.clone())]);
323 }
324 let text = std::fs::read_to_string(&path).unwrap();
325 assert!(text.len() <= LOG_MAX);
326 assert!(text.contains("https://699.test/") && !text.contains("https://0.test/"));
327 assert!(text.lines().all(|l| l.split('\t').count() == 4), "trimmed at a line boundary");
328 let _ = std::fs::remove_dir_all(dir);
329 }
330 use api::tests::{items, server};
331
332 fn scratch(name: &str) -> std::path::PathBuf {
333 let dir = std::env::temp_dir().join(format!("cce-raindrop-sync-{name}-{}", std::process::id()));
334 let _ = std::fs::remove_dir_all(&dir);
335 std::fs::create_dir_all(&dir).unwrap();
336 dir
337 }
338
339 #[test]
340 fn a_first_pass_imports_creates_and_saves_the_base() {
341 let dir = scratch("first");
342 let bookmarks = Bookmarks::at(dir.join("bookmarks.tsv"));
343 bookmarks.toggle("https://local.test/", "Mine");
344 // Raindrop: 1.test and 2.test; then the create for local.test.
345 let (base, seen) = server(vec![
346 (200, "", items(1..3, 2)),
347 (200, "", r#"{"result":true,"item":{"_id":90}}"#.into()),
348 ]);
349 let client = api::Client::with_base(Secret::from("t".to_string()), &base);
350 let base_path = dir.join("raindrop-sync.tsv");
351 let out = run_pass(&client, &bookmarks, &base_path, false).unwrap();
352 let Outcome::Synced(s) = out else { panic!("refused") };
353 assert_eq!((s.added_here, s.added_there), (2, 1));
354 assert_eq!(s.text(), "synced — 2 added here, 1 added to Raindrop");
355 let urls: Vec<_> = bookmarks.rows().into_iter().map(|r| r.1).collect();
356 assert_eq!(urls.len(), 3);
357 let saved = load_base(&base_path);
358 assert_eq!(saved.len(), 3, "both imports and the create are paired");
359 assert!(saved.iter().any(|b| b.id == 90 && b.url == "https://local.test/"));
360 assert_eq!(seen.lock().unwrap().len(), 2);
361 // The log names each bookmark that moved, and what happened to it.
362 let log = std::fs::read_to_string(log_path(&base_path)).unwrap();
363 let lines: Vec<Vec<&str>> = log.lines().map(|l| l.split('\t').collect()).collect();
364 assert_eq!(lines.len(), 3, "{log}");
365 assert!(lines.iter().all(|l| l.len() == 4 && l[0].ends_with('Z')));
366 assert!(lines.iter().any(|l| l[1] == "added here" && l[2] == "https://1.test/" && l[3] == "t1"));
367 assert!(lines.iter().any(|l| l[1] == "added to Raindrop" && l[2] == "https://local.test/" && l[3] == "Mine"));
368 let _ = std::fs::remove_dir_all(dir);
369 }
370
371 #[test]
372 fn a_refused_pass_changes_nothing_until_forced() {
373 let dir = scratch("refused");
374 let bookmarks = Bookmarks::at(dir.join("bookmarks.tsv"));
375 let base_path = dir.join("raindrop-sync.tsv");
376 // Synced before: three bookmarks, all gone here now.
377 save_base(&base_path, &(1..4).map(|i| super::super::Synced {
378 id: i, url: format!("https://{i}.test/"), title: format!("t{i}"),
379 }).collect::<Vec<_>>()).unwrap();
380 let (base, seen) = server(vec![
381 (200, "", items(1..4, 3)),
382 (200, "", items(1..4, 3)),
383 (200, "", "{}".into()),
384 (200, "", "{}".into()),
385 (200, "", "{}".into()),
386 ]);
387 let client = api::Client::with_base(Secret::from("t".to_string()), &base);
388 let out = run_pass(&client, &bookmarks, &base_path, false).unwrap();
389 assert!(matches!(out, Outcome::Refused(ref r) if r.contains("every bookmark in Raindrop")), "{out:?}");
390 assert_eq!(seen.lock().unwrap().len(), 1, "only the fetch went out");
391 assert_eq!(load_base(&base_path).len(), 3, "the base is untouched");
392 assert!(std::fs::read_to_string(log_path(&base_path)).unwrap().contains("\trefused\t\tthis would delete every bookmark in Raindrop"));
393
394 let Outcome::Synced(s) = run_pass(&client, &bookmarks, &base_path, true).unwrap() else {
395 panic!("forced pass refused")
396 };
397 assert_eq!(s.trashed_there, 3);
398 assert!(load_base(&base_path).is_empty());
399 let log = std::fs::read_to_string(log_path(&base_path)).unwrap();
400 assert_eq!(log.matches("\tmoved to Raindrop's trash\thttps://").count(), 3, "{log}");
401 let _ = std::fs::remove_dir_all(dir);
402 }
403
404 #[test]
405 fn an_idle_pass_does_not_touch_the_file() {
406 let dir = scratch("idle");
407 let path = dir.join("bookmarks.tsv");
408 let bookmarks = Bookmarks::at(path.clone());
409 let (base, _) = server(vec![(200, "", items(1..1, 0))]);
410 let client = api::Client::with_base(Secret::from("t".to_string()), &base);
411 let out = run_pass(&client, &bookmarks, &dir.join("raindrop-sync.tsv"), false).unwrap();
412 assert_eq!(out, Outcome::Synced(Summary::default()));
413 assert_eq!(Summary::default().text(), "in sync");
414 assert!(!path.exists(), "nothing to write, nothing written");
415 assert!(!log_path(&dir.join("raindrop-sync.tsv")).exists(), "an idle pass logs nothing");
416 let _ = std::fs::remove_dir_all(dir);
417 }
418 }