- lib.rs MAIN_WINDOW: the "main" window label that the reminder worker, the capture window and the shell each wrote, 5 sites in all. - lib.rs now_ms(): the epoch-milliseconds clock that the reminder worker and autosync's cycle stamp each computed. LastCycle.at_ms becomes i64, the same number on the wire. - integration.rs applications_dir(): the XDG launcher directory that the entry path and the desktop-database refresh each built. rustfmt --check is clean in the CI image. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
390 lines
14 KiB
Rust
390 lines
14 KiB
Rust
//! The desktop syncs on its own (#5167).
|
|
//!
|
|
//! Until this existed, the only thing that ever called the sync engine was the
|
|
//! "Sync now" button on the Sync screen. A linked desktop sat on its edits until
|
|
//! someone went looking for that button, and never heard about a note written on
|
|
//! the phone. Android has done better from the start (a periodic WorkManager job,
|
|
//! plus a flush when the app leaves the foreground); this is the desktop's version.
|
|
//!
|
|
//! ## When a cycle runs
|
|
//!
|
|
//! - **At launch**, so the board shows what changed while the app was closed.
|
|
//! - **Soon after an edit.** Every [`TICK`] the worker reads a fingerprint of what
|
|
//! is waiting to be sent (one cheap query), and sends it if that changed since the
|
|
//! last cycle. That puts an edit on the server within about ten seconds without a
|
|
//! hook in every command that writes, and a burst of autosaves still costs one
|
|
//! cycle per tick.
|
|
//!
|
|
//! The fingerprint, not "is anything pending", and that is the loop's fixed
|
|
//! point. A change the server rejects stays pending on purpose — it needs a
|
|
//! person — so "pending" would stay true after every cycle and the worker would
|
|
//! resend the same rejected row every ten seconds, forever, reporting success
|
|
//! each time. A rejected change still goes again with the next real edit, the
|
|
//! periodic cycle, or the button.
|
|
//! - **Every [`PULL_EVERY`]** even with nothing to send, so a quiet desktop still
|
|
//! receives what other devices wrote.
|
|
//! - **When the window regains focus**, which is when someone is about to look.
|
|
//! - **When asked**: the button, and the cycle right after linking.
|
|
//!
|
|
//! ## One cycle at a time, by construction
|
|
//!
|
|
//! Every cycle, the button's included, runs on this one worker thread. There is no
|
|
//! lock to forget: two cycles cannot overlap because there is only one place that
|
|
//! runs them.
|
|
//!
|
|
//! ## Failing quietly, but not silently
|
|
//!
|
|
//! An unreachable server is the normal state of a laptop on a train, so a failed
|
|
//! automatic cycle is not an error dialog. It backs off (doubling from [`TICK`] up
|
|
//! to [`MAX_BACKOFF`]), and its message is kept in [`LastCycle`] for the Sync
|
|
//! screen to show. The button ignores the backoff: a person pressing it wants to
|
|
//! know now.
|
|
//!
|
|
//! Not done: a final push when the app quits. Each request is bounded by the
|
|
//! client's timeout, but holding the window open for one on an offline laptop is a
|
|
//! worse quit than leaving the edit for the launch cycle, which sends it first.
|
|
|
|
use std::sync::mpsc::{self, Receiver, RecvTimeoutError, Sender};
|
|
use std::sync::Mutex;
|
|
use std::time::{Duration, Instant};
|
|
|
|
use serde::Serialize;
|
|
use tauri::{AppHandle, Emitter, Manager};
|
|
|
|
use inkwell_core::local::Db;
|
|
use inkwell_core::sync::blobs::BlobStore;
|
|
use inkwell_core::sync::engine::{self, SyncOutcome};
|
|
use inkwell_core::sync::{push, state};
|
|
|
|
/// How often the worker wakes to look for unsent edits.
|
|
const TICK: Duration = Duration::from_secs(10);
|
|
/// The longest a linked, idle desktop goes without asking the server for news.
|
|
const PULL_EVERY: Duration = Duration::from_secs(5 * 60);
|
|
/// The ceiling on the wait after repeated failures.
|
|
const MAX_BACKOFF: Duration = Duration::from_secs(5 * 60);
|
|
/// Focus comes and goes constantly; it may start a cycle at most this often.
|
|
const FOCUS_GAP: Duration = Duration::from_secs(30);
|
|
|
|
/// Emitted to every window after a cycle, with a [`LastCycle`] as the payload.
|
|
/// The board reloads when `changed` is true; the Sync screen redraws its status.
|
|
pub const SYNCED_EVENT: &str = "inkwell://synced";
|
|
|
|
/// What started a cycle. Shown on the Sync screen, so an automatic sync can be told
|
|
/// apart from a button press.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
|
|
#[serde(rename_all = "snake_case")]
|
|
pub enum Trigger {
|
|
Launch,
|
|
Edit,
|
|
Timer,
|
|
Focus,
|
|
Manual,
|
|
}
|
|
|
|
/// The outcome of the most recent cycle, automatic or not.
|
|
#[derive(Debug, Clone, Serialize)]
|
|
pub struct LastCycle {
|
|
/// When it finished, in milliseconds since the epoch — the frontend formats it.
|
|
pub at_ms: i64,
|
|
pub trigger: Trigger,
|
|
pub ok: bool,
|
|
/// Why it failed, in words a person can act on. None when `ok`.
|
|
pub error: Option<String>,
|
|
/// Whether the pull changed anything this device shows.
|
|
pub changed: bool,
|
|
/// The full result, for the screen that reports counts. None on failure.
|
|
pub outcome: Option<SyncOutcome>,
|
|
}
|
|
|
|
type Reply = Sender<Result<SyncOutcome, String>>;
|
|
|
|
struct Request {
|
|
trigger: Trigger,
|
|
/// Present when someone is waiting on the result (the button).
|
|
reply: Option<Reply>,
|
|
}
|
|
|
|
/// Managed state: the way into the worker, and the last thing it did.
|
|
pub struct AutoSync {
|
|
tx: Mutex<Sender<Request>>,
|
|
last: Mutex<Option<LastCycle>>,
|
|
}
|
|
|
|
impl AutoSync {
|
|
fn send(&self, request: Request) -> Result<(), String> {
|
|
self.tx
|
|
.lock()
|
|
.map_err(|e| e.to_string())?
|
|
.send(request)
|
|
.map_err(|_| "The sync worker has stopped.".to_string())
|
|
}
|
|
|
|
/// Ask for a cycle without waiting for it. Subject to the worker's backoff and,
|
|
/// for focus, to [`FOCUS_GAP`].
|
|
pub fn kick(&self, trigger: Trigger) {
|
|
if let Err(e) = self.send(Request {
|
|
trigger,
|
|
reply: None,
|
|
}) {
|
|
log::warn!("could not request a {trigger:?} sync: {e}");
|
|
}
|
|
}
|
|
|
|
pub fn last(&self) -> Option<LastCycle> {
|
|
self.last.lock().ok().and_then(|l| l.clone())
|
|
}
|
|
}
|
|
|
|
/// Start the worker and put [`AutoSync`] in managed state. Called from `setup`,
|
|
/// after the store and the blob store are managed — the worker reads both.
|
|
pub fn start(app: &AppHandle) -> std::io::Result<()> {
|
|
let (tx, rx) = mpsc::channel();
|
|
app.manage(AutoSync {
|
|
tx: Mutex::new(tx),
|
|
last: Mutex::new(None),
|
|
});
|
|
let handle = app.clone();
|
|
std::thread::Builder::new()
|
|
.name("autosync".into())
|
|
.spawn(move || worker(handle, rx))?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Run one cycle on the worker and wait for it — the "Sync now" button.
|
|
pub async fn run_now(app: &AppHandle) -> Result<SyncOutcome, String> {
|
|
let (reply, wait) = mpsc::channel();
|
|
app.state::<AutoSync>().send(Request {
|
|
trigger: Trigger::Manual,
|
|
reply: Some(reply),
|
|
})?;
|
|
// `recv` blocks, so it runs on the blocking pool rather than stalling the
|
|
// async runtime that every other command shares.
|
|
tauri::async_runtime::spawn_blocking(move || wait.recv())
|
|
.await
|
|
.map_err(|e| e.to_string())?
|
|
.map_err(|_| "The sync worker stopped before answering.".to_string())?
|
|
}
|
|
|
|
/// How long to wait after `failures` failed cycles in a row.
|
|
fn backoff(failures: u32) -> Duration {
|
|
if failures == 0 {
|
|
return Duration::ZERO;
|
|
}
|
|
TICK.saturating_mul(1 << failures.min(10)).min(MAX_BACKOFF)
|
|
}
|
|
|
|
/// The server URL and token, or None when this device isn't linked.
|
|
fn credentials(db: &Db) -> Result<Option<(String, String)>, String> {
|
|
let conn = db.conn()?;
|
|
state::credentials(&conn, None).map_err(|e| e.to_string())
|
|
}
|
|
|
|
/// What is waiting to be sent, as [`push::pending_fingerprint`] sees it. A store
|
|
/// that can't be read reads as nothing pending; the periodic cycle still runs.
|
|
fn pending(db: &Db) -> Option<String> {
|
|
let conn = db.conn().ok()?;
|
|
push::pending_fingerprint(&conn).ok().flatten()
|
|
}
|
|
|
|
/// Whether a cycle's pull changed anything the board shows.
|
|
fn pulled_changes(outcome: &SyncOutcome) -> bool {
|
|
let p = &outcome.pull;
|
|
p.notes_applied + p.notes_deleted + p.labels_applied + p.labels_deleted + p.blobs_downloaded > 0
|
|
}
|
|
|
|
/// The worker's whole state between cycles.
|
|
#[derive(Default)]
|
|
struct Pace {
|
|
last_attempt: Option<Instant>,
|
|
last_success: Option<Instant>,
|
|
failures: u32,
|
|
/// The pending fingerprint as the last cycle left it.
|
|
sent: Option<String>,
|
|
}
|
|
|
|
impl Pace {
|
|
/// Whether an automatic trigger may start a cycle now. The button never asks.
|
|
fn allows(&self, trigger: Trigger, now: Instant) -> bool {
|
|
let since_attempt = self.last_attempt.map(|t| now.duration_since(t));
|
|
// After a failure nothing automatic runs until the backoff has passed.
|
|
if since_attempt.is_some_and(|d| d < backoff(self.failures)) {
|
|
return false;
|
|
}
|
|
match trigger {
|
|
Trigger::Manual | Trigger::Launch | Trigger::Edit => true,
|
|
Trigger::Focus => since_attempt.is_none_or(|d| d >= FOCUS_GAP),
|
|
Trigger::Timer => self
|
|
.last_success
|
|
.is_none_or(|t| now.duration_since(t) >= PULL_EVERY),
|
|
}
|
|
}
|
|
}
|
|
|
|
fn worker(app: AppHandle, rx: Receiver<Request>) {
|
|
let mut pace = Pace::default();
|
|
let mut next = Some(Request {
|
|
trigger: Trigger::Launch,
|
|
reply: None,
|
|
});
|
|
loop {
|
|
let request = match next.take() {
|
|
Some(request) => Some(request),
|
|
None => match rx.recv_timeout(TICK) {
|
|
Ok(request) => Some(request),
|
|
Err(RecvTimeoutError::Timeout) => None,
|
|
// The app is shutting down: every sender went with the managed state.
|
|
Err(RecvTimeoutError::Disconnected) => return,
|
|
},
|
|
};
|
|
let db = app.state::<Db>();
|
|
|
|
let (trigger, reply) = match request {
|
|
Some(Request { trigger, reply }) => (trigger, reply),
|
|
// A quiet tick: new edits waiting to go, or the periodic pull.
|
|
None => match pending(&db) {
|
|
Some(now) if pace.sent.as_ref() != Some(&now) => (Trigger::Edit, None),
|
|
_ => (Trigger::Timer, None),
|
|
},
|
|
};
|
|
if reply.is_none() && !pace.allows(trigger, Instant::now()) {
|
|
continue;
|
|
}
|
|
|
|
let (base_url, token) = match credentials(&db) {
|
|
Ok(Some(creds)) => creds,
|
|
Ok(None) => {
|
|
// Unlinked is the normal resting state, not a failure.
|
|
if let Some(reply) = reply {
|
|
let _ = reply.send(Err("This app isn't linked to a server yet.".to_string()));
|
|
}
|
|
continue;
|
|
}
|
|
Err(e) => {
|
|
if let Some(reply) = reply {
|
|
let _ = reply.send(Err(e));
|
|
}
|
|
continue;
|
|
}
|
|
};
|
|
|
|
let blobs = app.state::<BlobStore>();
|
|
let started = Instant::now();
|
|
// No catch_unwind here, and that is deliberate. The workspace's release profile
|
|
// sets `panic = "abort"`, so a panic in a cycle ends the whole app rather than
|
|
// this thread, and a catch would only ever work in a debug build. What
|
|
// matters is that the worker can't die silently and leave the app looking
|
|
// synced, and an abort is anything but silent.
|
|
let result = tauri::async_runtime::block_on(engine::run_cycle(
|
|
db.inner(),
|
|
blobs.inner(),
|
|
&base_url,
|
|
&token,
|
|
));
|
|
pace.last_attempt = Some(started);
|
|
// Whatever is still pending after this cycle is not new on the next tick.
|
|
pace.sent = pending(&db);
|
|
match &result {
|
|
Ok(_) => {
|
|
pace.failures = 0;
|
|
pace.last_success = Some(started);
|
|
}
|
|
Err(e) => {
|
|
pace.failures = pace.failures.saturating_add(1);
|
|
log::warn!(
|
|
"{trigger:?} sync failed ({} in a row, next automatic try in {:?}): {e}",
|
|
pace.failures,
|
|
backoff(pace.failures)
|
|
);
|
|
}
|
|
}
|
|
|
|
let cycle = LastCycle {
|
|
at_ms: crate::now_ms(),
|
|
trigger,
|
|
ok: result.is_ok(),
|
|
error: result.as_ref().err().cloned(),
|
|
changed: result.as_ref().is_ok_and(pulled_changes),
|
|
outcome: result.as_ref().ok().cloned(),
|
|
};
|
|
if let Ok(mut last) = app.state::<AutoSync>().last.lock() {
|
|
*last = Some(cycle.clone());
|
|
}
|
|
if let Err(e) = app.emit(SYNCED_EVENT, &cycle) {
|
|
log::warn!("could not announce the finished sync: {e}");
|
|
}
|
|
if let Some(reply) = reply {
|
|
let _ = reply.send(result);
|
|
}
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn backoff_doubles_from_the_tick_and_stops_at_the_ceiling() {
|
|
assert_eq!(backoff(0), Duration::ZERO);
|
|
assert_eq!(backoff(1), TICK * 2);
|
|
assert_eq!(backoff(2), TICK * 4);
|
|
assert_eq!(backoff(40), MAX_BACKOFF);
|
|
}
|
|
|
|
#[test]
|
|
fn a_fresh_worker_runs_every_trigger() {
|
|
let pace = Pace::default();
|
|
let now = Instant::now();
|
|
for trigger in [
|
|
Trigger::Launch,
|
|
Trigger::Edit,
|
|
Trigger::Timer,
|
|
Trigger::Focus,
|
|
Trigger::Manual,
|
|
] {
|
|
assert!(pace.allows(trigger, now), "{trigger:?}");
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn after_a_failure_automatic_triggers_wait_out_the_backoff() {
|
|
let start = Instant::now();
|
|
let pace = Pace {
|
|
last_attempt: Some(start),
|
|
last_success: None,
|
|
failures: 1,
|
|
..Pace::default()
|
|
};
|
|
assert!(!pace.allows(Trigger::Edit, start + TICK));
|
|
assert!(!pace.allows(Trigger::Focus, start + TICK));
|
|
assert!(pace.allows(Trigger::Edit, start + backoff(1)));
|
|
}
|
|
|
|
#[test]
|
|
fn the_timer_only_pulls_once_the_interval_has_passed() {
|
|
let start = Instant::now();
|
|
let pace = Pace {
|
|
last_attempt: Some(start),
|
|
last_success: Some(start),
|
|
failures: 0,
|
|
..Pace::default()
|
|
};
|
|
assert!(!pace.allows(Trigger::Timer, start + TICK));
|
|
assert!(pace.allows(Trigger::Timer, start + PULL_EVERY));
|
|
// An edit does not wait for the interval.
|
|
assert!(pace.allows(Trigger::Edit, start + TICK));
|
|
}
|
|
|
|
#[test]
|
|
fn focus_is_rate_limited() {
|
|
let start = Instant::now();
|
|
let pace = Pace {
|
|
last_attempt: Some(start),
|
|
last_success: Some(start),
|
|
failures: 0,
|
|
..Pace::default()
|
|
};
|
|
assert!(!pace.allows(Trigger::Focus, start + Duration::from_secs(5)));
|
|
assert!(pace.allows(Trigger::Focus, start + FOCUS_GAP));
|
|
}
|
|
}
|