Skip to main content

confium_coordinator/
distributed_lock.rs

1//! Distributed lock manager — TTL-based leases with fencing tokens.
2
3use chrono::{DateTime, Duration, Utc};
4use serde::{Deserialize, Serialize};
5use std::collections::HashMap;
6use std::sync::Mutex;
7
8/// A distributed lock lease.
9#[derive(Debug, Clone, Serialize, Deserialize)]
10pub struct Lease {
11    pub resource: String,
12    pub holder: String,
13    pub fencing_token: u64,
14    pub acquired_at: DateTime<Utc>,
15    pub expires_at: DateTime<Utc>,
16}
17
18/// Thread-safe distributed lock manager (in-memory; production
19/// implementations use Redis or etcd).
20pub struct DistributedLockManager {
21    leases: Mutex<HashMap<String, Lease>>,
22    next_token: Mutex<u64>,
23}
24
25impl DistributedLockManager {
26    pub fn new() -> Self {
27        Self {
28            leases: Mutex::new(HashMap::new()),
29            next_token: Mutex::new(1),
30        }
31    }
32
33    /// Try to acquire a lock. Returns the fencing token if acquired,
34    /// or None if the resource is held by another holder.
35    pub fn try_acquire(&self, resource: &str, holder: &str, ttl: Duration) -> Option<u64> {
36        let mut leases = self.leases.lock().unwrap();
37        let now = Utc::now();
38
39        if let Some(existing) = leases.get(resource) {
40            if existing.holder != holder && existing.expires_at > now {
41                return None;
42            }
43        }
44
45        let mut next = self.next_token.lock().unwrap();
46        let token = *next;
47        *next += 1;
48
49        leases.insert(
50            resource.into(),
51            Lease {
52                resource: resource.into(),
53                holder: holder.into(),
54                fencing_token: token,
55                acquired_at: now,
56                expires_at: now + ttl,
57            },
58        );
59        Some(token)
60    }
61
62    /// Release a lock. Only the holder can release.
63    pub fn release(&self, resource: &str, holder: &str) -> bool {
64        let mut leases = self.leases.lock().unwrap();
65        if let Some(lease) = leases.get(resource) {
66            if lease.holder == holder {
67                leases.remove(resource);
68                return true;
69            }
70        }
71        false
72    }
73
74    /// Renew an existing lease. Returns the new expiry time.
75    pub fn renew(&self, resource: &str, holder: &str, ttl: Duration) -> Option<DateTime<Utc>> {
76        let mut leases = self.leases.lock().unwrap();
77        let lease = leases.get_mut(resource)?;
78        if lease.holder != holder {
79            return None;
80        }
81        lease.expires_at = Utc::now() + ttl;
82        Some(lease.expires_at)
83    }
84
85    /// Check if a resource is currently locked.
86    pub fn is_locked(&self, resource: &str) -> bool {
87        let leases = self.leases.lock().unwrap();
88        leases
89            .get(resource)
90            .map(|l| l.expires_at > Utc::now())
91            .unwrap_or(false)
92    }
93
94    /// Get the current holder of a resource.
95    pub fn holder(&self, resource: &str) -> Option<String> {
96        let leases = self.leases.lock().unwrap();
97        leases.get(resource).map(|l| l.holder.clone())
98    }
99
100    /// Purge expired leases.
101    pub fn purge_expired(&self) -> usize {
102        let mut leases = self.leases.lock().unwrap();
103        let now = Utc::now();
104        let before = leases.len();
105        leases.retain(|_, l| l.expires_at > now);
106        before - leases.len()
107    }
108
109    /// Number of active leases.
110    pub fn active_count(&self) -> usize {
111        let leases = self.leases.lock().unwrap();
112        let now = Utc::now();
113        leases.values().filter(|l| l.expires_at > now).count()
114    }
115}
116
117impl Default for DistributedLockManager {
118    fn default() -> Self {
119        Self::new()
120    }
121}
122
123#[cfg(test)]
124mod tests {
125    use super::*;
126
127    #[test]
128    fn acquire_returns_token() {
129        let mgr = DistributedLockManager::new();
130        let token = mgr.try_acquire("res", "a", Duration::seconds(30));
131        assert!(token.is_some());
132        assert!(token.unwrap() > 0);
133    }
134
135    #[test]
136    fn cannot_acquire_held_resource() {
137        let mgr = DistributedLockManager::new();
138        mgr.try_acquire("res", "a", Duration::seconds(30));
139        let token = mgr.try_acquire("res", "b", Duration::seconds(30));
140        assert!(token.is_none());
141    }
142
143    #[test]
144    fn same_holder_can_reacquire() {
145        let mgr = DistributedLockManager::new();
146        mgr.try_acquire("res", "a", Duration::seconds(30));
147        let token = mgr.try_acquire("res", "a", Duration::seconds(30));
148        assert!(token.is_some());
149    }
150
151    #[test]
152    fn release_by_holder() {
153        let mgr = DistributedLockManager::new();
154        mgr.try_acquire("res", "a", Duration::seconds(30));
155        assert!(mgr.release("res", "a"));
156        assert!(!mgr.is_locked("res"));
157    }
158
159    #[test]
160    fn release_by_wrong_holder_fails() {
161        let mgr = DistributedLockManager::new();
162        mgr.try_acquire("res", "a", Duration::seconds(30));
163        assert!(!mgr.release("res", "b"));
164    }
165
166    #[test]
167    fn renew_extends_lease() {
168        let mgr = DistributedLockManager::new();
169        mgr.try_acquire("res", "a", Duration::seconds(10));
170        let new_expiry = mgr.renew("res", "a", Duration::seconds(60));
171        assert!(new_expiry.is_some());
172    }
173
174    #[test]
175    fn renew_wrong_holder_fails() {
176        let mgr = DistributedLockManager::new();
177        mgr.try_acquire("res", "a", Duration::seconds(10));
178        assert!(mgr.renew("res", "b", Duration::seconds(60)).is_none());
179    }
180
181    #[test]
182    fn fencing_tokens_monotonic() {
183        let mgr = DistributedLockManager::new();
184        let t1 = mgr.try_acquire("r1", "a", Duration::seconds(30)).unwrap();
185        mgr.release("r1", "a");
186        let t2 = mgr.try_acquire("r2", "a", Duration::seconds(30)).unwrap();
187        assert!(t2 > t1);
188    }
189
190    #[test]
191    fn expired_lease_can_be_acquired() {
192        let mgr = DistributedLockManager::new();
193        mgr.try_acquire("res", "a", Duration::seconds(-1)); // already expired
194        assert!(!mgr.is_locked("res"));
195        let token = mgr.try_acquire("res", "b", Duration::seconds(30));
196        assert!(token.is_some());
197    }
198
199    #[test]
200    fn purge_expired_removes_stale() {
201        let mgr = DistributedLockManager::new();
202        mgr.try_acquire("r1", "a", Duration::seconds(-1)); // expired
203        mgr.try_acquire("r2", "b", Duration::seconds(30));
204        assert_eq!(mgr.purge_expired(), 1);
205        assert_eq!(mgr.active_count(), 1);
206    }
207
208    #[test]
209    fn holder_returns_current() {
210        let mgr = DistributedLockManager::new();
211        mgr.try_acquire("res", "alice", Duration::seconds(30));
212        assert_eq!(mgr.holder("res").as_deref(), Some("alice"));
213    }
214}