From a36236efff49d6524c0b79dee582b44a59b01ccf Mon Sep 17 00:00:00 2001 From: John Coffey Date: Sun, 27 Sep 2026 16:26:04 -0700 Subject: [PATCH] End a locked account's delegation at its date A delegation with an end date dropped out of the delegate's token then, but its folder grants stayed until the daily sweep, so the delegate kept the account as an ordinary share for up to a day. Each node now sleeps until the soonest end date, woken early by any lock write and at least hourly, and re-applies that lock under a cluster-wide claim. The sweep also had a second-run bug: a delegation past its date gave the delegate back its earlier share, then dropped the note, so the next sweep removed that share entirely. The note is now kept while the delegate is still listed. --- crates/features/src/lock/mod.rs | 93 ++++++++++++++++++++++- crates/services/src/inbuxa_lock_expiry.rs | 71 +++++++++++++++++ crates/services/src/lib.rs | 5 ++ tests/src/system/account_lock.rs | 33 ++++++++ 4 files changed, 199 insertions(+), 3 deletions(-) create mode 100644 crates/services/src/inbuxa_lock_expiry.rs diff --git a/crates/features/src/lock/mod.rs b/crates/features/src/lock/mod.rs index b19e4ff..d44a26d 100644 --- a/crates/features/src/lock/mod.rs +++ b/crates/features/src/lock/mod.rs @@ -27,6 +27,32 @@ use store::{ write::{AnyClass, BatchBuilder, ValueClass}, }; use trc::AddContext; + +/// Rung when a lock is written, so this node's expiry timer re-reads the +/// `until` dates (AL-5): a delegation ends at its time, not at a sweep. +pub static UNTIL_CHANGED: tokio::sync::Notify = tokio::sync::Notify::const_new(); + +/// The soonest `until` still ahead of `now`, across every lock. +pub fn next_until(locks: &[Lock], now: u64) -> Option { + locks + .iter() + .flat_map(|lock| &lock.delegates) + .filter_map(|delegate| delegate.until) + .filter(|until| *until > now) + .min() +} + +/// Locks with a delegation that ended in `(after, now]`. +pub fn ended_between(locks: &[Lock], after: u64, now: u64) -> impl Iterator + '_ { + locks + .iter() + .filter(move |lock| { + lock.delegates + .iter() + .any(|d| d.until.is_some_and(|until| until > after && until <= now)) + }) + .map(|lock| lock.account_id) +} use types::{ acl::{Acl, AclGrant}, collection::Collection, @@ -240,10 +266,19 @@ pub fn merge_grants( if is_current(new, delegate.account_id) { continue; } - let before = noted(old, delegate.account_id) + let note = noted(old, delegate.account_id); + let before = note + .as_ref() .map(|r| Bitmap::from(r.rights)) .unwrap_or_default(); set(&mut acls, delegate.account_id, before); + // Still listed but past its `until`: keep the note, so running + // this again puts back the same share instead of removing it + if let Some(note) = note + && new.is_some_and(|new| new.delegate(delegate.account_id).is_some()) + { + replaced.push(note); + } } } @@ -415,8 +450,9 @@ pub async fn set(data: &Store, lock: &Lock, previous: Option<&Lock>) -> trc::Res batch.set(class(KIND_LOCK, &[lock.account_id]), Json(lock).serialize()?); data.write(batch.build_all()) .await - .caused_by(trc::location!()) - .map(|_| ()) + .caused_by(trc::location!())?; + UNTIL_CHANGED.notify_one(); + Ok(()) } /// Removes a lock and its delegate index. @@ -532,6 +568,57 @@ mod tests { assert!(after.contains(&AclGrant { account_id: 3, grants: read })); } + #[test] + fn an_expired_delegation_gives_back_its_share_every_time() { + let earlier: Bitmap = Bitmap::from_iter([Acl::Read]); + let note = Replaced { + collection: Collection::Mailbox as u8, + document_id: 5, + delegate: 2, + rights: u64::from(earlier), + }; + let mut ending = delegate(2, Access::Full); + ending.until = Some(200); + let lock = lock_with(vec![ending], vec![note.clone()]); + let during = vec![AclGrant { + account_id: 2, + grants: Access::Full.grants(Collection::Mailbox, false), + }]; + + // At its `until`, the share it had before comes back, and the note stays + let mut replaced = Vec::new(); + let after = merge_grants(&during, Collection::Mailbox, 5, false, Some(&lock), Some(&lock), 300, &mut replaced) + .unwrap(); + assert_eq!(after, vec![AclGrant { account_id: 2, grants: earlier }]); + assert_eq!(replaced, vec![note.clone()]); + + // The next sweep changes nothing, rather than removing that share + let swept = Lock { replaced: replaced.clone(), ..lock }; + let mut again = Vec::new(); + assert!( + merge_grants(&after, Collection::Mailbox, 5, false, Some(&swept), Some(&swept), 400, &mut again).is_none() + ); + assert_eq!(again, vec![note]); + } + + #[test] + fn the_timer_finds_the_next_end() { + let ends_at = |account_id, until| { + let mut d = delegate(account_id, Access::Read); + d.until = until; + d + }; + let a = Lock { account_id: 10, ..lock_with(vec![ends_at(2, Some(500)), ends_at(3, None)], vec![]) }; + let b = Lock { account_id: 11, ..lock_with(vec![ends_at(4, Some(300))], vec![]) }; + let locks = vec![a, b]; + assert_eq!(next_until(&locks, 100), Some(300)); + assert_eq!(next_until(&locks, 300), Some(500)); + assert_eq!(next_until(&locks, 500), None); + assert_eq!(ended_between(&locks, 100, 300).collect::>(), vec![11]); + assert_eq!(ended_between(&locks, 300, 600).collect::>(), vec![10]); + assert!(ended_between(&locks, 600, 900).next().is_none()); + } + #[test] fn expired_delegations_grant_nothing() { let lock = Lock { diff --git a/crates/services/src/inbuxa_lock_expiry.rs b/crates/services/src/inbuxa_lock_expiry.rs new file mode 100644 index 0000000..00f2b27 --- /dev/null +++ b/crates/services/src/inbuxa_lock_expiry.rs @@ -0,0 +1,71 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! Ends a locked account's delegations at their `until` (AL-5), rather than +//! at the next daily sweep. Each node sleeps until the soonest `until`, wakes +//! early when a lock is written here, and checks at least hourly for locks +//! written on other nodes. One node re-applies each lock; the others find it +//! claimed. + +use common::{BuildServer, Inner, KV_LOCK_TASK, Server}; +use inbuxa_features::lock; +use std::{sync::Arc, time::Duration}; +use store::write::now; + +/// The longest the timer sleeps, so an `until` set on another node is seen. +const CEILING: u64 = 3600; + +pub fn spawn_lock_expiry(inner: Arc) { + tokio::spawn(async move { + // Since boot: anything that ended while the server was down + let mut checked = 0; + loop { + let server = inner.build_server(); + let now = now(); + let wait = match lock::all(server.store()).await { + Ok(locks) => { + for account_id in lock::ended_between(&locks, checked, now).collect::>() + { + end_delegations(&server, account_id).await; + } + checked = now; + lock::next_until(&locks, now) + .map_or(CEILING, |until| (until - now).min(CEILING)) + } + Err(err) => { + trc::error!(err.details("Failed to read account locks for their end dates")); + 60 + } + }; + tokio::select! { + _ = tokio::time::sleep(Duration::from_secs(wait.max(1))) => {} + _ = lock::UNTIL_CHANGED.notified() => {} + } + } + }); +} + +async fn end_delegations(server: &Server, account_id: u32) { + let key = [b"lock-until:".as_slice(), &account_id.to_be_bytes()].concat(); + match server + .in_memory_store() + .try_lock(KV_LOCK_TASK, &key, 60) + .await + { + Ok(true) => { + if let Err(err) = email::inbuxa_lock::reconcile(server, account_id).await { + trc::error!( + err.account_id(account_id) + .details("Failed to end a delegation at its date") + ); + } + } + Ok(false) => {} + Err(err) => { + trc::error!(err.details("Failed to claim a delegation's end")); + } + } +} diff --git a/crates/services/src/lib.rs b/crates/services/src/lib.rs index 682247f..c1e76d9 100644 --- a/crates/services/src/lib.rs +++ b/crates/services/src/lib.rs @@ -23,6 +23,8 @@ use std::sync::Arc; use crate::task_manager::{manager::spawn_task_manager, scheduler::spawn_task_scheduler}; pub mod broadcast; +// inbuxa: AL-5, delegations end at their date +pub mod inbuxa_lock_expiry; pub mod state_manager; pub mod task_manager; @@ -65,6 +67,9 @@ impl SpawnServices for IpcReceivers { // Spawn task manager spawn_task_manager(inner.clone()); + // inbuxa: AL-5, end delegations at their `until` + inbuxa_lock_expiry::spawn_lock_expiry(inner.clone()); + // Spawn task scheduler spawn_task_scheduler(inner); } diff --git a/tests/src/system/account_lock.rs b/tests/src/system/account_lock.rs index 316dd15..40067d0 100644 --- a/tests/src/system/account_lock.rs +++ b/tests/src/system/account_lock.rs @@ -313,6 +313,39 @@ pub async fn test(test: &mut TestServer) { "AL-4: the rejected message wasn't kept: {kept}" ); + // AL-5: a delegation ends at its `until`, not at the next daily sweep + let soon = store::write::now() + 3; + let response = admin + .lock_set(json!({"reason": "Handover ends shortly", + "update": {owner_id.as_str(): {"delegates": [ + {"accountId": delegate.id_string(), "access": "organize", + "until": chrono::DateTime::from_timestamp(soon as i64, 0).unwrap().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)}]}}})) + .await; + assert!( + response["updated"].get(owner_id.as_str()).is_some(), + "AL-5: {response}" + ); + let (_, before) = delegate + .call("Mailbox/get", json!({"accountId": owner_id, "ids": null})) + .await; + assert!( + before["list"].as_array().is_some_and(|l| !l.is_empty()), + "AL-5: the delegate lost the account before its end: {before}" + ); + tokio::time::sleep(Duration::from_secs(6)).await; + let session = delegate.jmap_session_object().await.0; + assert!( + session["accounts"].get(owner_id.as_str()).is_none(), + "AL-5: the delegation outlived its end in the session: {session}" + ); + let (_, after) = delegate + .call("Mailbox/get", json!({"accountId": owner_id, "ids": null})) + .await; + assert!( + after["list"].as_array().is_none_or(|l| l.is_empty()), + "AL-5: the delegate still reaches the folders after its end: {after}" + ); + // AL-10: unlocking needs a reason, then restores everything let response = admin .lock_set(json!({"destroy": [owner_id]})) -- 2.54.0