diff --git a/Cargo.lock b/Cargo.lock index 10a23c0..e820995 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -19,6 +19,18 @@ dependencies = [ "cpufeatures 0.2.17", ] +[[package]] +name = "ahash" +version = "0.8.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a15f179cd60c4584b8a8c596927aadc462e27f2ca70c04e0071964a73ba7a75" +dependencies = [ + "cfg-if", + "once_cell", + "version_check", + "zerocopy", +] + [[package]] name = "ascii" version = "1.1.0" @@ -556,6 +568,18 @@ dependencies = [ "pin-project-lite", ] +[[package]] +name = "fallible-iterator" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2acce4a10f12dc2fb14a218589d4f1f62ef011b2d0cc4b3cb1bba8e94da14649" + +[[package]] +name = "fallible-streaming-iterator" +version = "0.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7360491ce676a36bf9bb3c56c1aa791658183a54d2744120f27285738d90465a" + [[package]] name = "fastrand" version = "2.5.0" @@ -956,12 +980,30 @@ dependencies = [ "system-deps", ] +[[package]] +name = "hashbrown" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" +dependencies = [ + "ahash", +] + [[package]] name = "hashbrown" version = "0.17.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ed5909b6e89a2db4456e54cd5f673791d7eca6732202bbf2a9cc504fe2f9b84a" +[[package]] +name = "hashlink" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6ba4ff7128dee98c7dc9794b6a411377e1404dba1c97deb8d1a55297bd25d8af" +dependencies = [ + "hashbrown 0.14.5", +] + [[package]] name = "heck" version = "0.5.0" @@ -1036,7 +1078,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d466e9454f08e4a911e14806c24e16fba1b4c121d1ea474396f396069cf949d9" dependencies = [ "equivalent", - "hashbrown", + "hashbrown 0.17.1", ] [[package]] @@ -1165,6 +1207,17 @@ version = "0.2.189" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2" +[[package]] +name = "libsqlite3-sys" +version = "0.30.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e99fb7a497b1e3339bc746195567ed8d3e24945ecd636e3619d20b9de9e9149" +dependencies = [ + "cc", + "pkg-config", + "vcpkg", +] + [[package]] name = "linux-raw-sys" version = "0.12.1" @@ -1243,6 +1296,7 @@ dependencies = [ "notify-rust", "rand", "roxmltree", + "rusqlite", "rustls", "rustls-native-certs", "secret-service", @@ -1593,6 +1647,20 @@ dependencies = [ "memchr", ] +[[package]] +name = "rusqlite" +version = "0.32.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7753b721174eb8ff87a9a0e799e2d7bc3749323e773db92e0984debb00019d6e" +dependencies = [ + "bitflags", + "fallible-iterator", + "fallible-streaming-iterator", + "hashlink", + "libsqlite3-sys", + "smallvec", +] + [[package]] name = "rustc_version" version = "0.4.1" @@ -2190,6 +2258,12 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "vcpkg" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" + [[package]] name = "version-compare" version = "0.2.1" @@ -2634,6 +2708,26 @@ dependencies = [ "serde", ] +[[package]] +name = "zerocopy" +version = "0.8.56" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "556764e583adb45a9f8d413c2a147fa7e8d821e48e12b14fd560b607998b75eb" +dependencies = [ + "zerocopy-derive", +] + +[[package]] +name = "zerocopy-derive" +version = "0.8.56" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2ab42fc20575779bd240faa45f94a74256f755c0fa9e89f0ede20d91d0cdfc1" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "zeroize" version = "1.9.0" diff --git a/Cargo.toml b/Cargo.toml index cdb8934..ddc6325 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -18,6 +18,7 @@ libadwaita = { version = "0.9", features = ["v1_5"] } notify = "8" rand = "0.10" roxmltree = "0.21" +rusqlite = { version = "0.32", features = ["bundled"] } rustls = "0.23" rustls-native-certs = "0.8" secret-service = { version = "5", features = ["rt-tokio-crypto-rust"] } diff --git a/src/core/account_runtime.rs b/src/core/account_runtime.rs index 63541a7..bcb949c 100644 --- a/src/core/account_runtime.rs +++ b/src/core/account_runtime.rs @@ -167,12 +167,19 @@ impl WatcherTargets { } } - /// A remote file hint triggers a remote-push sync on every folder, like the - /// Python `NotifyPushClient` callback `scheduler.request(REMOTE_PUSH)`. - fn apply_remote_push(targets: &Rc>) { + /// Route a remote file hint to the folders that contain the notified files. + /// + /// With `notify_file_id` (issue #183) each notified id maps to the folder + /// whose external journal knows that file (`has_file_ids`), so only those + /// folders get a remote-push sync instead of re-syncing every folder of + /// the account. A legacy `notify_file` (empty ids) fans out to every + /// folder, preserving the original `NotifyPushClient` behaviour. + fn apply_remote_push(targets: &Rc>, file_ids: Vec) { let schedulers = targets.borrow().schedulers.clone(); for scheduler in &schedulers { - scheduler.request(Trigger::RemotePush); + if file_ids.is_empty() || scheduler.has_file_ids(&file_ids) { + scheduler.request(Trigger::RemotePush); + } } } @@ -180,6 +187,13 @@ impl WatcherTargets { fn store_push_state(targets: &Rc>, state: PushState, message: String) { if let Ok(mut current) = targets.try_borrow_mut() { current.push_state = Some((state, message)); + // Issue #185: while push is connected it is the real-time source of + // truth, so the periodic remote-interval poll can be skipped. Tell + // every folder scheduler whether to drop `RemoteInterval`. + let ready = state == PushState::Connected; + for scheduler in ¤t.schedulers { + scheduler.set_remote_push_ready(ready); + } } } @@ -876,7 +890,7 @@ impl AccountRuntime { let state_targets = Rc::clone(&self.targets); Some(NotifyPushClient::new( self.account.provider, - move || WatcherTargets::apply_remote_push(&file_targets), + move |file_ids| WatcherTargets::apply_remote_push(&file_targets, file_ids), move || WatcherTargets::apply_server_notification(¬ification_targets), move |state, message| WatcherTargets::store_push_state(&state_targets, state, message), )) @@ -1213,9 +1227,11 @@ impl AccountRuntime { } /// Drive the notify_push file-notification fan-out (what the client's - /// `on_file_notification` callback runs in production). - pub(crate) fn simulate_remote_push(&self) { - WatcherTargets::apply_remote_push(&self.targets); + /// `on_file_notification` callback runs in production). An empty list is + /// the legacy `notify_file` fan-out; non-empty ids route to the folders + /// that contain them (issue #183). + pub(crate) fn simulate_remote_push(&self, file_ids: Vec) { + WatcherTargets::apply_remote_push(&self.targets, file_ids); } /// Build the push client for the current account (test mirror of @@ -1933,7 +1949,7 @@ mod tests { // `simulate_remote_push` runs exactly what the NotifyPushClient // `on_file_notification` callback runs in production. - runtime.simulate_remote_push(); + runtime.simulate_remote_push(Vec::new()); assert!(source.borrow().pending() >= 1); } diff --git a/src/core/files_journal.rs b/src/core/files_journal.rs new file mode 100644 index 0000000..ac05f5b --- /dev/null +++ b/src/core/files_journal.rs @@ -0,0 +1,165 @@ +//! Read-only access to the external sync engine journal. +//! +//! `nextcloudcmd`/`opencloudcmd` keep their own reconciliation journal per +//! local folder: a `nextcloud`-style SQLite database named `.sync_.db` +//! inside the folder's local root. The journal stores one row per known file, +//! including its remote file id (`metadata.fileid`). +//! +//! The official Nextcloud desktop client resolves `notify_file_id` push hints +//! against exactly this table via `SyncJournalDb::hasFileIds`, comparing the +//! numeric file ids from the push hint against the `metadata.fileid` column +//! cast to an integer (the column holds a string like `00152532ocwv4xsuk6ni`, +//! whose numeric prefix is the file id). This module mirrors that lookup so a +//! remote change can be routed to the folder that actually contains the +//! notified file, instead of re-syncing every folder of the account (issue +//! #183). +//! +//! It is intentionally read-only and best-effort: a missing journal (first +//! sync, or the folder was never touched) simply answers `false`, which is the +//! conservative fallback for the push routing. + +use std::path::{Path, PathBuf}; + +/// Find the sync engine journal for a local folder root. +/// +/// The journal is a `.sync_.db` file whose name is derived from the +/// folder path and account, so we match any `.sync*.db` regular file directly +/// under `local_root` (the engines never nest them). Returns the first match, +/// or `None` when the folder has never been reconciled. +pub fn find_journal(local_root: &Path) -> Option { + let entries = std::fs::read_dir(local_root).ok()?; + for entry in entries.flatten() { + let path = entry.path(); + let file_name = entry.file_name(); + let name = file_name.to_string_lossy(); + if name.starts_with(".sync_") && name.ends_with(".db") && path.is_file() { + return Some(path); + } + } + None +} + +/// Whether any of `file_ids` is present in the journal for `local_root`. +/// +/// Opens the journal read-only (without connecting to a database managed +/// concurrently by the engine) and answers whether at least one notified id +/// maps to a file this folder knows about. Mirrors +/// `SyncJournalDb::hasFileIds`: the `metadata.fileid` string column is +/// compared to each given integer id after a `CAST(... AS INTEGER)`. +/// +/// Returns `false` when the journal cannot be read (missing, busy, or the +/// query fails), which routes the hint conservatively. +pub fn contains_file_ids(local_root: &Path, file_ids: &[i64]) -> bool { + if file_ids.is_empty() { + return false; + } + let Some(journal) = find_journal(local_root) else { + return false; + }; + journal_contains(&journal, file_ids) +} + +/// Open `journal` and test membership of any `file_ids` in `metadata`. +fn journal_contains(journal: &Path, file_ids: &[i64]) -> bool { + let conn = match rusqlite::Connection::open_with_flags( + journal, + rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_URI, + ) { + Ok(conn) => conn, + Err(_) => return false, + }; + // The journal may be mid-write by the engine; a busy query is not a + // mismatch, so fall back to `false` (the caller routes conservatively). + let sql = match build_file_id_membership_sql(file_ids.len()) { + Ok(sql) => sql, + Err(_) => return false, + }; + let mut stmt = match conn.prepare(&sql) { + Ok(stmt) => stmt, + Err(_) => return false, + }; + let params: Vec = file_ids.to_vec(); + stmt.query_row(rusqlite::params_from_iter(¶ms), |_| Ok(())) + .is_ok() +} + +/// Build the membership query for `n` ids: a single placeholder per id +/// (`CAST(fileid AS INTEGER) IN (?1, ?2, ...)`). The `metadata.fileid` +/// column holds a string like `00152532ocwv4xsuk6ni`, so we cast it to an +/// integer to match the numeric ids sent by the push hint (the official +/// client compares the same way). +fn build_file_id_membership_sql(count: usize) -> Result { + if count == 0 { + return Err("no file ids to test".to_string()); + } + let mut sql = String::from("SELECT 1 FROM metadata WHERE CAST(fileid AS INTEGER) IN ("); + for i in 0..count { + if i > 0 { + sql.push(','); + } + sql.push('?'); + } + sql.push_str(") LIMIT 1"); + Ok(sql) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn write_journal(db_path: &Path, rows: &[(&str, &str)]) { + let conn = rusqlite::Connection::open(db_path).unwrap(); + conn.execute_batch( + "CREATE TABLE metadata(phash INTEGER(8),pathlen INTEGER,path VARCHAR(4096), + inode INTEGER,uid INTEGER,gid INTEGER,mode INTEGER,modtime INTEGER(8), + type INTEGER,md5 VARCHAR(32), fileid VARCHAR(128), remotePerm VARCHAR(128), + filesize BIGINT, ignoredChildrenRemote INT);", + ) + .unwrap(); + let mut stmt = conn + .prepare("INSERT INTO metadata(path, fileid) VALUES (?1, ?2)") + .unwrap(); + for (path, fileid) in rows { + stmt.execute([path, fileid]).unwrap(); + } + } + + #[test] + fn find_journal_returns_sync_db_in_folder_root() { + let dir = tempfile::tempdir().unwrap(); + let db = dir.path().join(".sync_abc123.db"); + write_journal(&db, &[("/file.txt", "152532ocwv4xsuk6ni")]); + assert_eq!(find_journal(dir.path()), Some(db)); + } + + #[test] + fn contains_file_ids_matches_numeric_prefix_of_string_id() { + let dir = tempfile::tempdir().unwrap(); + write_journal( + &dir.path().join(".sync_x.db"), + &[("/a/b.txt", "00152532ocwv4")], + ); + assert!(contains_file_ids(dir.path(), &[152532])); + assert!(!contains_file_ids(dir.path(), &[999])); + } + + #[test] + fn contains_file_ids_matches_any_of_multiple_ids() { + let dir = tempfile::tempdir().unwrap(); + write_journal(&dir.path().join(".sync_x.db"), &[("/a/b.txt", "17ocwv4")]); + assert!(contains_file_ids(dir.path(), &[5, 17])); + } + + #[test] + fn missing_journal_answers_false() { + let dir = tempfile::tempdir().unwrap(); + assert!(!contains_file_ids(dir.path(), &[42])); + } + + #[test] + fn empty_id_list_answers_false() { + let dir = tempfile::tempdir().unwrap(); + write_journal(&dir.path().join(".sync_x.db"), &[("/a", "42")]); + assert!(!contains_file_ids(dir.path(), &[])); + } +} diff --git a/src/core/mod.rs b/src/core/mod.rs index 2022478..35d8533 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -8,6 +8,7 @@ pub mod debounce; pub mod delete_guard; pub mod desktop_integration; pub mod exclusions; +pub mod files_journal; pub mod log; pub mod network; pub mod notifications; diff --git a/src/core/scheduler.rs b/src/core/scheduler.rs index 7743594..08f42e6 100644 --- a/src/core/scheduler.rs +++ b/src/core/scheduler.rs @@ -46,6 +46,20 @@ pub const KEYRING_RETRY_BASE_MS: u64 = 2000; /// without hammering the network. pub const SERVER_PROBE_INTERVAL_MS: u64 = 30_000; +/// Return the delay (seconds) before the next run after `consecutive` +/// failing syncs, mirroring the official client's `scheduleFolder` backoff +/// (issue #184): a folder that fails repeatedly is not re-run on every +/// trigger in a tight loop (3-4 fails -> 10s, 5-6 -> 30s, more -> 60s). +/// Runs with fewer than three consecutive failures are scheduled normally. +pub fn failure_backoff_seconds(consecutive: u32) -> u64 { + match consecutive { + 0..=2 => 0, + 3..=4 => 10, + 5..=6 => 30, + _ => 60, + } +} + /// How a finished reconciliation turned out. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum SyncOutcome { @@ -135,6 +149,18 @@ struct SchedulerInner { /// locked collection must not retry forever: once the budget is spent the /// folder parks in the keyring-locked state until a manual action. keyring_retry_count: u32, + /// Issue #184: how many consecutive failed syncs this folder has had. + /// Mirrors the official client's `consecutiveFailingSyncs`: a run that + /// fails backs off (10/30/60s) instead of re-running on every trigger in a + /// tight loop. Reset on any successful/conflicted run. + consecutive_failing_syncs: u32, + /// Issue #185: whether this account's notify_push channel is delivering + /// change notifications (see [`Scheduler::set_remote_push_ready`]). When + /// it is, the periodic remote-interval poll is dropped: push is the + /// real-time source of truth and the interval would only re-run a full + /// reconciliation on top of it (the official client skips the ETag poll + /// for accounts whose push is ready). + remote_push_ready: bool, delete_alert: Option, delete_bypass_once: bool, /// Whether a synchronization has ever completed successfully. Drives the @@ -224,6 +250,8 @@ impl Scheduler { auth_required: false, server_unreachable: false, run_active: None, + consecutive_failing_syncs: 0, + remote_push_ready: false, }; let inner = Rc::new(RefCell::new(inner)); { @@ -311,6 +339,27 @@ impl Scheduler { self.inner.borrow_mut().local_root = local_root; } + /// Tell the scheduler whether the account's notify_push channel is ready + /// (issue #185). When it is, periodic remote-interval polls are dropped: + /// push already surfaces remote changes in real time, and the interval + /// would only run a full reconciliation on top of them. + pub fn set_remote_push_ready(&self, ready: bool) { + self.inner.borrow_mut().remote_push_ready = ready; + } + + /// Whether this folder's external sync journal knows any of `file_ids`. + /// + /// Used to route a `notify_file_id` push hint to exactly the folder that + /// contains the notified file, instead of re-syncing every folder of the + /// account (issue #183). Returns `false` when the folder has no local + /// root yet or the journal is unavailable; the caller then falls back to + /// the conservative fan-out. + pub fn has_file_ids(&self, file_ids: &[i64]) -> bool { + let root = self.inner.borrow().local_root.clone(); + root.map(|root| crate::core::files_journal::contains_file_ids(&root, file_ids)) + .unwrap_or(false) + } + /// Share the run-active flag with the progress forwarder (issue #145): /// the forwarder must not repaint the label from stale events once a run /// has finished and cleared it. @@ -456,6 +505,13 @@ impl SchedulerInner { if self.stopped { return; } + // Issue #185: while the account's push channel is ready, a periodic + // remote-interval poll is redundant (push delivers changes in real + // time) and would only force a full reconciliation. Drop it; the remote + // Interval timer stays armed and is skipped until push is unready. + if trigger == Trigger::RemoteInterval && self.remote_push_ready { + return; + } if self.delete_alert.is_some() && !self.delete_bypass_once { self.queue.add(trigger); let message = self @@ -622,6 +678,36 @@ impl SchedulerInner { self.start_source = Some(id); } + /// Seconds to wait before re-running after repeated failures (issue + /// #184), mirroring the official client's delay from + /// `consecutiveFailingSyncs`. Zero means no backoff (schedule now). + fn failure_backoff(&self) -> u64 { + failure_backoff_seconds(self.consecutive_failing_syncs) + } + + /// Re-arm `start` after `seconds`, using the same timeout source as the + /// server probe and keyring retry. Guards against double-arming. + fn schedule_after(&mut self, seconds: u64) { + if self.stopped + || self.start_source.is_some() + || self.preparing + || self.running + || seconds == 0 + { + return; + } + let weak = self.self_ref.clone(); + let id = self.source.borrow_mut().add_timeout( + Duration::from_secs(seconds), + Box::new(move || { + if let Some(inner) = weak.upgrade() { + inner.borrow_mut().start(); + } + }), + ); + self.start_source = Some(id); + } + fn start(&mut self) { self.start_source = None; if self.stopped || self.preparing || self.running || self.queue.is_empty() || !self.online { @@ -777,6 +863,7 @@ impl SchedulerInner { self.keyring_retry_count = 0; self.auth_required = false; self.server_unreachable = false; + self.consecutive_failing_syncs = 0; self.ever_synced = true; self.set_idle_state(); (true, false) @@ -786,6 +873,7 @@ impl SchedulerInner { self.keyring_retry_count = 0; self.auth_required = false; self.server_unreachable = false; + self.consecutive_failing_syncs = 0; self.ever_synced = true; self.state.set( AppState::IdleOk, @@ -798,6 +886,7 @@ impl SchedulerInner { // Issue #72: arm the credential gate so the periodic and // inotify triggers stop retrying a dead password. self.auth_required = true; + self.consecutive_failing_syncs += 1; self.state.set( AppState::AuthRequired, t("Credentials rejected. Sign in again in Account settings."), @@ -809,6 +898,7 @@ impl SchedulerInner { // No stored credentials: park the folder like a rejection so // the triggers do not hammer the server (issue #95). self.auth_required = true; + self.consecutive_failing_syncs += 1; self.state.set( AppState::AuthRequired, t("No saved credentials. Use Sign in again in Account settings."), @@ -817,12 +907,14 @@ impl SchedulerInner { } SyncOutcome::KeyringLocked => { self.keyring_locked = true; + self.consecutive_failing_syncs += 1; self.state .set(AppState::KeyringLocked, t("Password keyring is locked")); (false, false) } SyncOutcome::Failed => { self.keyring_locked = false; + self.consecutive_failing_syncs += 1; self.state .set(AppState::Error, t("Synchronization failed — view the log")); (true, false) @@ -837,6 +929,7 @@ impl SchedulerInner { // account/folder, never to other accounts. self.keyring_locked = false; self.server_unreachable = true; + self.consecutive_failing_syncs += 1; self.state.set( AppState::Offline, t("Synchronization blocked: the server is unreachable"), @@ -931,7 +1024,15 @@ impl SchedulerInner { // (schedule_start aborts on in_cooldown) and must start now // (issue #127). if !self.queue.is_empty() && self.online && !self.paused() { - self.schedule_start(); + let delay = self.failure_backoff(); + if delay == 0 { + self.schedule_start(); + } else { + // Issue #184: a folder that keeps failing backs off instead of + // re-running on every trigger in a tight loop (this mirrors the + // official client's delay tied to consecutiveFailingSyncs). + self.schedule_after(delay); + } } else if !self.paused() { self.set_idle_state(); } @@ -2410,4 +2511,36 @@ mod tests { // Degenerate values never match. assert!(!time_in_window("12:00", "bad", "06:30")); } + + /// Issue #184: the backoff table mirrors the official client's + /// `scheduleFolder` delays (10/30/60s) based on consecutive failures. + #[test] + fn failure_backoff_seconds_matches_the_terminal_delays() { + assert_eq!(failure_backoff_seconds(0), 0); + assert_eq!(failure_backoff_seconds(2), 0); + assert_eq!(failure_backoff_seconds(3), 10); + assert_eq!(failure_backoff_seconds(4), 10); + assert_eq!(failure_backoff_seconds(5), 30); + assert_eq!(failure_backoff_seconds(6), 30); + assert_eq!(failure_backoff_seconds(7), 60); + assert_eq!(failure_backoff_seconds(100), 60); + } + + /// Issue #185: when the account's push channel is ready, a periodic + /// remote-interval poll is dropped instead of queuing a full sync. + #[test] + fn remote_interval_is_dropped_while_push_is_ready() { + let (scheduler, source, _runner) = make_scheduler(None); + scheduler.set_remote_push_ready(true); + scheduler.request(Trigger::RemoteInterval); + assert!(source.borrow().pending() == 0, "no start was scheduled"); + assert_eq!(scheduler.queue_len(), 0, "the interval was not queued"); + // Disabling push readiness lets the interval through again. + scheduler.set_remote_push_ready(false); + scheduler.request(Trigger::RemoteInterval); + assert!( + source.borrow().pending() >= 1, + "interval now schedules a run" + ); + } } diff --git a/src/nextcloud/push.rs b/src/nextcloud/push.rs index e1f5135..7edbbfc 100644 --- a/src/nextcloud/push.rs +++ b/src/nextcloud/push.rs @@ -32,6 +32,7 @@ use std::rc::Rc; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; use std::time::Duration; +use std::time::Instant; use rustls::pki_types::ServerName; use rustls::{ClientConfig, ClientConnection, RootCertStore, StreamOwned}; @@ -53,6 +54,20 @@ type PushStream = MaybeTlsStream; /// Reconnect delays in seconds, mirroring `NotifyPushClient.BACKOFF_SECONDS`. pub const BACKOFF_SECONDS: [u64; 6] = [2, 5, 10, 30, 60, 300]; +/// How long (ms) the push worker tolerates total connection inactivity +/// (no inbound text/ping/pong) before declaring the WebSocket dead and +/// reconnecting (issue #186). The official client keeps a 30s ping/pong +/// heartbeat; a half-dead TCP connection (e.g. a stuck proxy) otherwise sits +/// silently with no notifications flowing and no reconnect. 75s leaves room +/// for a normal idle server and a slow pong while still catching zombies. +pub const PUSH_HEARTBEAT_TIMEOUT_MS: u64 = 75_000; + +/// Maximum consecutive failed push authentication attempts before the client +/// stops retrying on its own and parks in `AuthRequired` (issue #186), +/// mirroring the official client's `MAX_ALLOWED_FAILED_AUTHENTICATION_ATTEMPTS` +/// (3). Prevents an invalid credential from retrying forever on a timer. +pub const MAX_PUSH_AUTH_ATTEMPTS: u32 = 3; + /// The backoff delay for a given failed-attempt index (capped at the last /// value, exactly like `BACKOFF_SECONDS[min(index, len - 1)]` in Python). pub fn backoff_seconds(backoff_index: usize) -> u64 { @@ -92,8 +107,10 @@ enum PushEventKind { Unsupported, /// The WebSocket sent `authenticated`. Authenticated, - /// The WebSocket sent a `notify_file*` hint. - FileNotification, + /// The WebSocket sent a `notify_file*` hint. Carries the numeric file ids + /// from a `notify_file_id ` hint (issue #183); empty for a + /// legacy `notify_file` hint (no ids), meaning an unknown location. + FileNotification(Vec), /// The WebSocket sent a `notify_notification` hint (new server /// notification, issue #31). Notification, @@ -122,11 +139,15 @@ struct PushInner { generation: u64, backoff_index: usize, force_password_auth: bool, + /// Issue #186: consecutive failed authentication attempts. Mirrors the + /// official client's cap (3); once exceeded the client parks in + /// `AuthRequired` instead of reconnecting in a loop on a bad credential. + auth_failure_count: u32, connection: Option, reconnect: Option>, tls: Option>, backoff_scale: f64, - on_file_notification: Rc, + on_file_notification: Rc)>, on_notification: Rc, on_state: Rc, } @@ -143,13 +164,15 @@ pub struct NotifyPushClient { impl NotifyPushClient { /// Create a client for an account bound to `provider`. /// - /// `on_file_notification` fires on every remote file hint, `on_notification` + /// `on_file_notification` fires on every remote file hint, carrying the + /// notified file ids (empty for a legacy `notify_file` hint), + /// `on_notification` /// on every `notify_notification` server hint, `on_state` on every /// [`PushState`] transition (message included). All three run on the main /// thread. pub fn new( provider: Provider, - on_file_notification: impl Fn() + 'static, + on_file_notification: impl Fn(Vec) + 'static, on_notification: impl Fn() + 'static, on_state: impl Fn(PushState, String) + 'static, ) -> Self { @@ -164,6 +187,7 @@ impl NotifyPushClient { generation: 0, backoff_index: 0, force_password_auth: false, + auth_failure_count: 0, connection: None, reconnect: None, tls: None, @@ -213,7 +237,16 @@ impl NotifyPushClient { /// Reflect the network status; disconnects (keeping the config) when /// going offline and reconnects when back online. pub fn set_online(&self, online: bool) { - self.inner.borrow_mut().online = online; + { + let mut inner = self.inner.borrow_mut(); + inner.online = online; + if online { + // Issue #186: coming back online (or an explicit re-enable) is + // a fresh start — clear the auth-failure streak so a channel + // parked after repeated auth failures can retry. + inner.auth_failure_count = 0; + } + } if !online { self.disconnect(true); } else if self.inner.borrow().enabled { @@ -334,12 +367,13 @@ impl NotifyPushClient { let mut inner = self.inner.borrow_mut(); inner.force_password_auth = false; inner.backoff_index = 0; + inner.auth_failure_count = 0; drop(inner); self.on_state(PushState::Connected, "Connected"); } - PushEventKind::FileNotification => { + PushEventKind::FileNotification(file_ids) => { let callback = self.inner.borrow().on_file_notification.clone(); - callback(); + callback(file_ids); } PushEventKind::Notification => { let callback = self.inner.borrow().on_notification.clone(); @@ -368,7 +402,25 @@ impl NotifyPushClient { Some(details) } }; - if let Some(details) = details { + // Issue #186: a connection that keeps closing before it + // authenticates is burning credentials on a loop. Mirror the + // official client's cap: after MAX_PUSH_AUTH_ATTEMPTS failed + // attempts, stop reconnecting and park the channel in + // AuthRequired until an explicit action (network restored / + // settings change) restarts it. + let exhausted = { + let mut inner = self.inner.borrow_mut(); + if !authenticated { + inner.auth_failure_count += 1; + } + inner.auth_failure_count >= MAX_PUSH_AUTH_ATTEMPTS + }; + if exhausted { + self.on_state( + PushState::AuthRequired, + "Push authentication failed after several attempts.", + ); + } else if let Some(details) = details { self.schedule_reconnect(details); } } @@ -584,27 +636,62 @@ fn push_worker_main(inputs: WorkerInputs) { } let mut authenticated = false; + // Issue #186: track the last inbound activity (any text/ping/pong) so a + // half-dead connection (no pongs, no server pings) is detected and + // torn down; the main side reconnects through the existing backoff. + let mut last_activity = Instant::now(); loop { if stop.load(Ordering::SeqCst) { graceful_close(&mut websocket); return; } + // Heartbeat: if nothing has arrived for the timeout, the connection is + // dead (or the server stopped sending pings). Close so the main side + // reconnects rather than silently missing all notifications. + if last_activity.elapsed().as_millis() as u64 > PUSH_HEARTBEAT_TIMEOUT_MS { + let _ = tx.send_blocking(PushEvent { + generation, + kind: PushEventKind::Closed { + reason: "Push heartbeat timed out: no activity from the server.".to_string(), + auth_mode, + authenticated, + }, + }); + graceful_close(&mut websocket); + return; + } match websocket.read() { Ok(Message::Text(text)) => { + last_activity = Instant::now(); let text = text.as_str().trim(); if text == "authenticated" { authenticated = true; + // Issue #183: opt in to the modern file-id notifications so + // a change can be routed to the exact folder (the official + // client sends `listen n_id` right after authenticating). + let _ = websocket.send(Message::text("listen n_id")); let _ = tx.send_blocking(PushEvent { generation, kind: PushEventKind::Authenticated, }); - } else if text == "notify_file" - || text == "notify_file_id" - || text.starts_with("notify_file_id ") - { + } else if text == "notify_file" { let _ = tx.send_blocking(PushEvent { generation, - kind: PushEventKind::FileNotification, + kind: PushEventKind::FileNotification(Vec::new()), + }); + } else if text == "notify_file_id" || text.starts_with("notify_file_id ") { + // Issue #183: the modern hint carries a JSON array of + // numeric file ids (`notify_file_id [141, 142]`). Parse it + // so the event can be routed to the folder that contains + // the file; an empty or unparsable list falls back to the + // legacy fan-out (empty ids). + let ids = text + .strip_prefix("notify_file_id") + .map(parse_file_id_list) + .unwrap_or_default(); + let _ = tx.send_blocking(PushEvent { + generation, + kind: PushEventKind::FileNotification(ids), }); } else if text == "notify_notification" { let _ = tx.send_blocking(PushEvent { @@ -619,9 +706,12 @@ fn push_worker_main(inputs: WorkerInputs) { } } Ok(Message::Ping(_)) => { + last_activity = Instant::now(); let _ = websocket.flush(); } - Ok(Message::Pong(_)) | Ok(Message::Frame(_)) | Ok(Message::Binary(_)) => {} + Ok(Message::Pong(_)) | Ok(Message::Frame(_)) | Ok(Message::Binary(_)) => { + last_activity = Instant::now(); + } Ok(Message::Close(_)) => break, Err(tungstenite::Error::ConnectionClosed) => break, Err(tungstenite::Error::Io(error)) if is_would_block(&error) => continue, @@ -698,6 +788,25 @@ fn parse_capabilities(body: &[u8]) -> Result, String> { parse_push_capability(&payload) } +/// Parse the id list trailing a `notify_file_id` hint. +/// +/// The server sends `notify_file_id [141, 142]` (a JSON array of numeric +/// file ids). Anything that is not a JSON array of integers yields an empty +/// list, which the caller treats as the legacy "unknown file" fan-out. +fn parse_file_id_list(trailing: &str) -> Vec { + let trimmed = trailing.trim(); + if !trimmed.starts_with('[') { + return Vec::new(); + } + let Ok(value) = serde_json::from_str::(trimmed) else { + return Vec::new(); + }; + let Some(array) = value.as_array() else { + return Vec::new(); + }; + array.iter().filter_map(|entry| entry.as_i64()).collect() +} + /// Extract the pre-auth token from the response body (plain text or JSON). /// /// Mirrors `_request_pre_auth`: a JSON body yields `token` or @@ -1073,7 +1182,7 @@ mod tests { let notifications_clone = Rc::clone(¬ifications); let client = NotifyPushClient::new( provider, - move || { + move |_file_ids| { notifications_clone.set(notifications_clone.get() + 1); }, || {}, @@ -1409,7 +1518,7 @@ mod tests { .iter() .map(|event| match &event.kind { PushEventKind::Authenticated => "authenticated", - PushEventKind::FileNotification => "file_notification", + PushEventKind::FileNotification(_) => "file_notification", PushEventKind::Closed { .. } => "closed", other => panic!("unexpected event: {other:?}"), }) @@ -1445,7 +1554,7 @@ mod tests { .iter() .map(|event| match &event.kind { PushEventKind::Authenticated => "authenticated", - PushEventKind::FileNotification => "file_notification", + PushEventKind::FileNotification(_) => "file_notification", PushEventKind::Notification => "notification", PushEventKind::Closed { .. } => "closed", other => panic!("unexpected event: {other:?}"), @@ -1688,4 +1797,20 @@ mod tests { }) .expect("the test main context is available"); } + + #[test] + fn parse_file_id_list_reads_a_json_array_of_ints() { + assert_eq!(parse_file_id_list(" [141, 142]"), vec![141, 142]); + } + + #[test] + fn parse_file_id_list_returns_empty_for_non_array_or_non_int() { + assert!(parse_file_id_list("").is_empty()); + assert!(parse_file_id_list("not a list").is_empty()); + assert!(parse_file_id_list("{}").is_empty()); + // Non-numeric entries are dropped but numeric ones are kept (the + // official client ignores non-integer array entries). + assert_eq!(parse_file_id_list("[141, \"two\"]"), vec![141]); + assert!(parse_file_id_list("[\"a\",\"b\"]").is_empty()); + } }