End a locked account's delegation at its date #66

Merged
jcoffey-dev merged 1 commits from fix/delegation-until into main 2026-09-27 23:31:30 +00:00
4 changed files with 199 additions and 3 deletions
Showing only changes of commit a36236efff - Show all commits
+90 -3
View File
@@ -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<u64> {
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<Item = u32> + '_ {
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<Acl> = 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<_>>(), vec![11]);
assert_eq!(ended_between(&locks, 300, 600).collect::<Vec<_>>(), vec![10]);
assert!(ended_between(&locks, 600, 900).next().is_none());
}
#[test]
fn expired_delegations_grant_nothing() {
let lock = Lock {
+71
View File
@@ -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<Inner>) {
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::<Vec<_>>()
{
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"));
}
}
}
+5
View File
@@ -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);
}
+33
View File
@@ -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]}))