desktop: syncs on its own — at launch, soon after an edit, every few minutes and on focus
Android / Build, or is the channel already serving this? (push) Successful in 5s
CI & Build / Build now, or wait for Android? (push) Successful in 2s
CI & Build / Python lint (push) Successful in 2s
CI & Build / Web typecheck and unit tests (push) Successful in 13s
CI & Build / Python tests (push) Successful in 14s
Desktop (Tauri) / Build, or is the channel already serving this? (push) Successful in 2s
CI & Build / integration (push) Successful in 47s
CI & Build / Build & push image (push) Skipped
Desktop (Tauri) / Web tests, clippy, Rust tests and rustfmt (push) Successful in 3m8s
Desktop (Tauri) / Windows installer (cross-compiled) (push) Successful in 2m53s
Desktop (Tauri) / Tauri desktop (Linux) (push) Successful in 4m47s
Desktop (Tauri) / Update manifest (push) Successful in 6s
Android / Kotlin + Rust (APK) (push) Successful in 11m34s

Until now the only caller of the sync engine was the "Sync now" button. A worker
thread now owns every cycle (the button's included, so two never overlap):

- launch: one cycle as the app opens;
- edit: every 10s it reads a fingerprint of the pending set and sends when that
  moved. A fingerprint rather than "anything pending", because a rejected change
  stays pending and would otherwise be resent every tick forever;
- timer: a pull every 5 minutes with nothing to send;
- focus: at most once per 30s.

Failed automatic cycles back off (doubling from 10s to 5 minutes). A panicking
cycle counts as a failed one rather than ending the thread. Every cycle is
emitted as inkwell://synced: the board reloads when the pull changed something,
and the Sync screen shows the last automatic failure. No final push on quit;
the launch cycle sends whatever was left.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-10-06 23:40:27 -04:00
co-authored by Claude Opus 5.5
parent 09239ee3c7
commit efb141e555
8 changed files with 586 additions and 43 deletions
+393
View File
@@ -0,0 +1,393 @@
//! 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, SystemTime, UNIX_EPOCH};
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: u64,
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.0.lock().map_err(|e| e.to_string())?;
let current = state::read(&conn).map_err(|e| e.to_string())?;
Ok(current.server_url.zip(current.device_token))
}
/// 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.0.lock().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();
// A panic inside one cycle must not end automatic sync for the rest of the
// session: this thread is the only thing that runs it, and a dead thread
// looks exactly like a quiet one. It becomes a failed cycle instead.
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
tauri::async_runtime::block_on(engine::run_cycle(
db.inner(),
blobs.inner(),
&base_url,
&token,
))
}))
.unwrap_or_else(|_| Err("The sync cycle crashed; it will be retried.".to_string()));
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: SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |d| d.as_millis() as u64),
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));
}
}
+15 -19
View File
@@ -4,16 +4,17 @@
//! any of this. The Settings UI (M10.7e) drives these.
use serde::{Deserialize, Serialize};
use tauri::State;
use tauri::{AppHandle, State};
use inkwell_core::local::Db;
use inkwell_core::sync::blobs::BlobStore;
use inkwell_core::sync::client::{self, Identity, ProbeResult};
use inkwell_core::sync::compat::Compatibility;
use inkwell_core::sync::engine;
use inkwell_core::sync::push;
use inkwell_core::sync::state;
use crate::autosync::{self, AutoSync, LastCycle};
/// Ask a server who it is, without committing to anything. The UI calls this as the
/// user finishes typing an address, so they see what answered before handing over
/// credentials.
@@ -156,29 +157,24 @@ pub fn sync_status(db: State<'_, Db>) -> Result<state::Status, String> {
state::status(&conn).map_err(|e| e.to_string())
}
/// The server URL + token, or a plain "not linked" error. Every networked sync
/// command needs exactly this, and none of them may hold the lock past it.
fn credentials(db: &State<'_, Db>) -> Result<(String, String), String> {
let conn = db.0.lock().map_err(|e| e.to_string())?;
let current = state::read(&conn).map_err(|e| e.to_string())?;
match (current.server_url, current.device_token) {
(Some(url), Some(token)) => Ok((url, token)),
_ => Err("This app isn't linked to a server yet.".to_string()),
}
}
/// Run one full sync: push local changes, then pull the server's.
///
/// The only sync entry point exposed to the UI, on purpose. Push and pull exist
/// separately inside the crate, but offering a bare "pull" would let the UI overwrite
/// unsent local edits — the ordering isn't a suggestion, it's what keeps them.
///
/// Runs on the background worker (autosync.rs) and waits for it, so a button press
/// can never overlap an automatic cycle.
#[tauri::command]
pub async fn sync_now(
db: State<'_, Db>,
blobs: State<'_, BlobStore>,
) -> Result<engine::SyncOutcome, String> {
let (base_url, token) = credentials(&db)?;
engine::run_cycle(db.inner(), blobs.inner(), &base_url, &token).await
pub async fn sync_now(app: AppHandle) -> Result<engine::SyncOutcome, String> {
autosync::run_now(&app).await
}
/// The most recent cycle, automatic or not — what the Sync screen shows now that
/// most cycles happen without anyone pressing anything. None until the first one.
#[tauri::command]
pub fn sync_last(autosync: State<'_, AutoSync>) -> Option<LastCycle> {
autosync.last()
}
/// Whether anything is waiting to be sent. Lets the UI show an honest "unsynced
+17
View File
@@ -9,6 +9,7 @@
//! remains here is the Tauri command surface (`commands`), desktop integration
//! (menu-entry install for the Linux AppImage), the in-app updater, and boot.
mod autosync;
mod capture;
mod commands;
mod crossover;
@@ -142,6 +143,21 @@ pub fn run() {
// was built before this path could be resolved.
sync::blobs::publish_root(blobs.root().to_path_buf());
app.manage(blobs);
// Last: the worker reads the store and the blob store, both managed now.
// Its first cycle is the launch sync.
autosync::start(app.handle())?;
// Coming back to the window is when someone is about to look, so it asks
// for a cycle. Rate-limited in autosync, because focus flaps constantly.
if let Some(window) = app.get_webview_window("main") {
let handle = app.handle().clone();
window.on_window_event(move |event| {
if let tauri::WindowEvent::Focused(true) = event {
handle
.state::<autosync::AutoSync>()
.kick(autosync::Trigger::Focus);
}
});
}
Ok(())
})
.invoke_handler(tauri::generate_handler![
@@ -186,6 +202,7 @@ pub fn run() {
commands::sync::sync_unlink,
commands::sync::sync_status,
commands::sync::sync_now,
commands::sync::sync_last,
commands::sync::sync_has_pending,
update::update_channel_get,
update::update_channel_set,