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]}))