confium_coordinator/
distributed_lock.rs1use chrono::{DateTime, Duration, Utc};
4use serde::{Deserialize, Serialize};
5use std::collections::HashMap;
6use std::sync::Mutex;
7
8#[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
18pub 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 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 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 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 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 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 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 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)); 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)); 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}