git.lucas.co / cce-browser
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 }