1use crate::shared_lease_state::{
14 read_epoch, read_state, reject_symlink, remove_state, write_epoch, write_state,
15};
16use crate::{ProviderError, ProviderResult};
17use fs2::FileExt;
18use std::fs::{self, OpenOptions};
19use std::path::{Path, PathBuf};
20
21#[derive(Debug, Clone, PartialEq, Eq, Hash)]
23pub struct LeaseOwner(String);
24
25impl LeaseOwner {
26 pub fn new(value: impl Into<String>) -> ProviderResult<Self> {
28 let value = value.into();
29 validate_label("lease owner", &value)?;
30 Ok(Self(value))
31 }
32
33 pub fn as_str(&self) -> &str {
35 &self.0
36 }
37}
38
39#[derive(Debug, Clone, PartialEq, Eq)]
41pub struct LeaseToken {
42 pub resource: String,
44 pub owner: LeaseOwner,
46 pub epoch: u64,
48}
49
50#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct SharedResourceLease {
53 pub token: LeaseToken,
55 pub acquired_at_ms: u64,
57 pub heartbeat_at_ms: u64,
59 pub expires_at_ms: u64,
61}
62
63impl SharedResourceLease {
64 pub fn is_expired(&self, now_ms: u64, policy: &LeasePolicy) -> bool {
66 now_ms.saturating_sub(policy.clock_skew_ms) >= self.expires_at_ms
67 }
68}
69
70#[derive(Debug, Clone, PartialEq, Eq)]
72pub struct LeaseHeartbeat {
73 pub token: LeaseToken,
75 pub now_ms: u64,
77}
78
79#[derive(Debug, Clone, Copy, PartialEq, Eq)]
81pub struct LeasePolicy {
82 pub ttl_ms: u64,
84 pub heartbeat_ms: u64,
86 pub clock_skew_ms: u64,
88}
89
90impl LeasePolicy {
91 pub fn new(ttl_ms: u64, heartbeat_ms: u64, clock_skew_ms: u64) -> ProviderResult<Self> {
93 let policy = Self {
94 ttl_ms,
95 heartbeat_ms,
96 clock_skew_ms,
97 };
98 policy.validate()?;
99 Ok(policy)
100 }
101
102 fn validate(self) -> ProviderResult<()> {
103 if self.ttl_ms == 0
104 || self.heartbeat_ms == 0
105 || self.ttl_ms <= self.heartbeat_ms
106 || self.clock_skew_ms >= self.ttl_ms
107 {
108 return Err(ProviderError::InvalidConfiguration(
109 "lease ttl must exceed heartbeat and clock skew".to_string(),
110 ));
111 }
112 Ok(())
113 }
114}
115
116#[derive(Debug, Clone, Copy, PartialEq, Eq)]
118pub enum LeaseDecision {
119 Allowed,
121 NoLease,
123 Expired,
125 WrongOwner,
127 StaleEpoch,
129}
130
131pub trait LeaseRepository: Send + Sync {
133 fn acquire(
135 &self,
136 resource: &str,
137 owner: LeaseOwner,
138 now_ms: u64,
139 ) -> ProviderResult<SharedResourceLease>;
140 fn heartbeat(&self, heartbeat: LeaseHeartbeat) -> ProviderResult<SharedResourceLease>;
142 fn release(&self, token: &LeaseToken) -> ProviderResult<()>;
144 fn current(&self, resource: &str) -> ProviderResult<Option<SharedResourceLease>>;
146 fn check_fence(&self, token: &LeaseToken, now_ms: u64) -> ProviderResult<LeaseDecision>;
148}
149
150#[derive(Debug, Clone)]
152pub struct FileLeaseRepository {
153 root: PathBuf,
154 policy: LeasePolicy,
155}
156
157impl FileLeaseRepository {
158 pub fn open(root: impl Into<PathBuf>, policy: LeasePolicy) -> ProviderResult<Self> {
160 let root = root.into();
161 policy.validate()?;
162 reject_symlink(&root)?;
163 fs::create_dir_all(&root)
164 .map_err(|error| initialization(format!("lease root create failed: {error}")))?;
165 reject_symlink(&root)?;
166 if !root.is_dir() {
167 return Err(invalid("lease root is not a directory"));
168 }
169 Ok(Self { root, policy })
170 }
171
172 pub fn policy(&self) -> LeasePolicy {
174 self.policy
175 }
176
177 fn lock_path(&self, resource: &str) -> ProviderResult<PathBuf> {
178 Ok(self.root.join(format!("{}.lock", resource_file(resource)?)))
179 }
180
181 fn state_path(&self, resource: &str) -> ProviderResult<PathBuf> {
182 Ok(self
183 .root
184 .join(format!("{}.lease", resource_file(resource)?)))
185 }
186
187 fn epoch_path(&self, resource: &str) -> ProviderResult<PathBuf> {
188 Ok(self
189 .root
190 .join(format!("{}.epoch", resource_file(resource)?)))
191 }
192
193 fn with_lock<T>(
194 &self,
195 resource: &str,
196 operation: impl FnOnce(&Path) -> ProviderResult<T>,
197 ) -> ProviderResult<T> {
198 reject_symlink(&self.root)?;
199 let lock_path = self.lock_path(resource)?;
200 reject_symlink(&lock_path)?;
201 let file = OpenOptions::new()
202 .create(true)
203 .truncate(false)
204 .read(true)
205 .write(true)
206 .open(lock_path)
207 .map_err(|error| initialization(format!("lease lock open failed: {error}")))?;
208 file.lock_exclusive()
209 .map_err(|error| initialization(format!("lease lock failed: {error}")))?;
210 let result = operation(&self.state_path(resource)?);
211 FileExt::unlock(&file).map_err(|_| initialization("lease unlock failed"))?;
212 result
213 }
214}
215
216impl LeaseRepository for FileLeaseRepository {
217 fn acquire(
218 &self,
219 resource: &str,
220 owner: LeaseOwner,
221 now_ms: u64,
222 ) -> ProviderResult<SharedResourceLease> {
223 self.with_lock(resource, |state_path| {
224 let current = read_state(state_path, resource)?;
225 if let Some(current) = ¤t {
226 if !current.is_expired(now_ms, &self.policy) && current.token.owner != owner {
227 return Err(ProviderError::Initialization(
228 "shared resource lease is held by another owner".to_string(),
229 ));
230 }
231 }
232 let epoch_path = self.epoch_path(resource)?;
233 let previous_epoch = read_epoch(&epoch_path)?
234 .into_iter()
235 .chain(current.iter().map(|lease| lease.token.epoch))
236 .max()
237 .unwrap_or(0);
238 let next_epoch = previous_epoch
239 .checked_add(1)
240 .ok_or_else(|| initialization("shared resource fencing epoch exhausted"))?;
241 let expires_at_ms = lease_expiration(now_ms, self.policy.ttl_ms)?;
242 let lease = SharedResourceLease {
243 token: LeaseToken {
244 resource: resource.to_string(),
245 owner,
246 epoch: next_epoch,
247 },
248 acquired_at_ms: now_ms,
249 heartbeat_at_ms: now_ms,
250 expires_at_ms,
251 };
252 write_epoch(&epoch_path, next_epoch)?;
253 write_state(state_path, &lease)?;
254 Ok(lease)
255 })
256 }
257
258 fn heartbeat(&self, heartbeat: LeaseHeartbeat) -> ProviderResult<SharedResourceLease> {
259 self.with_lock(&heartbeat.token.resource, |state_path| {
260 let mut current =
261 read_state(state_path, &heartbeat.token.resource)?.ok_or_else(|| {
262 ProviderError::Initialization("shared resource lease is absent".to_string())
263 })?;
264 ensure_current(¤t, &heartbeat.token, heartbeat.now_ms, self.policy)?;
265 if heartbeat.now_ms < current.heartbeat_at_ms {
266 return Err(initialization("lease heartbeat clock moved backwards"));
267 }
268 current.heartbeat_at_ms = heartbeat.now_ms;
269 current.expires_at_ms = lease_expiration(heartbeat.now_ms, self.policy.ttl_ms)?;
270 write_state(state_path, ¤t)?;
271 Ok(current)
272 })
273 }
274
275 fn release(&self, token: &LeaseToken) -> ProviderResult<()> {
276 self.with_lock(&token.resource, |state_path| {
277 let current = read_state(state_path, &token.resource)?.ok_or_else(|| {
278 ProviderError::Initialization("shared resource lease is absent".to_string())
279 })?;
280 ensure_same_token(¤t, token)?;
281 remove_state(state_path)
282 })
283 }
284
285 fn current(&self, resource: &str) -> ProviderResult<Option<SharedResourceLease>> {
286 self.with_lock(resource, |state_path| read_state(state_path, resource))
287 }
288
289 fn check_fence(&self, token: &LeaseToken, now_ms: u64) -> ProviderResult<LeaseDecision> {
290 self.with_lock(&token.resource, |state_path| {
291 let Some(current) = read_state(state_path, &token.resource)? else {
292 return Ok(LeaseDecision::NoLease);
293 };
294 if current.is_expired(now_ms, &self.policy) {
295 return Ok(LeaseDecision::Expired);
296 }
297 if current.token.owner != token.owner {
298 return Ok(LeaseDecision::WrongOwner);
299 }
300 if current.token.epoch != token.epoch {
301 return Ok(LeaseDecision::StaleEpoch);
302 }
303 Ok(LeaseDecision::Allowed)
304 })
305 }
306}
307
308fn lease_expiration(now_ms: u64, ttl_ms: u64) -> ProviderResult<u64> {
309 now_ms
310 .checked_add(ttl_ms)
311 .ok_or_else(|| initialization("lease expiration timestamp overflow"))
312}
313
314fn ensure_current(
315 current: &SharedResourceLease,
316 token: &LeaseToken,
317 now_ms: u64,
318 policy: LeasePolicy,
319) -> ProviderResult<()> {
320 ensure_same_token(current, token)?;
321 if current.is_expired(now_ms, &policy) {
322 return Err(ProviderError::Initialization(
323 "shared resource lease has expired".to_string(),
324 ));
325 }
326 Ok(())
327}
328
329fn ensure_same_token(current: &SharedResourceLease, token: &LeaseToken) -> ProviderResult<()> {
330 if current.token != *token {
331 return Err(ProviderError::Initialization(
332 "shared resource lease fencing token mismatch".to_string(),
333 ));
334 }
335 Ok(())
336}
337
338fn resource_file(resource: &str) -> ProviderResult<String> {
339 validate_label("lease resource", resource)?;
340 Ok(resource.to_string())
341}
342
343fn validate_label(field: &'static str, value: &str) -> ProviderResult<()> {
344 if value.is_empty()
345 || value.len() > 128
346 || !value
347 .bytes()
348 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'-' | b'_'))
349 {
350 return Err(ProviderError::InvalidConfiguration(format!(
351 "invalid {field}"
352 )));
353 }
354 Ok(())
355}
356
357fn initialization(message: impl Into<String>) -> ProviderError {
358 ProviderError::Initialization(message.into())
359}
360
361fn invalid(message: impl Into<String>) -> ProviderError {
362 ProviderError::InvalidConfiguration(message.into())
363}
364
365#[cfg(test)]
366#[path = "shared_lease_tests.rs"]
367mod tests;