use std::fs::{File, OpenOptions, TryLockError};
use std::io::{Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use parking_lot::Mutex;
pub const MAINTENANCE_LOCK_FILE: &str = "maintenance.lock";
const PID_GATE_SUFFIX: &str = ".gate";
const PID_GATE_WAIT: Duration = Duration::from_millis(500);
const PID_GATE_POLL: Duration = Duration::from_millis(1);
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum LeaseStatus {
Held,
HeldElsewhere { holder_pid: Option<u32> },
Unavailable { reason: String },
}
impl LeaseStatus {
#[must_use]
pub fn is_held(&self) -> bool {
matches!(self, Self::Held)
}
}
#[derive(Debug)]
pub struct MaintenanceLease {
path: PathBuf,
gate_wait: Duration,
inner: Mutex<Inner>,
}
#[derive(Debug, Default)]
struct Inner {
file: Option<File>,
last_denial: Option<LeaseStatus>,
}
impl MaintenanceLease {
#[must_use]
pub fn new(data_root: &Path) -> Self {
Self {
path: data_root.join(MAINTENANCE_LOCK_FILE),
gate_wait: PID_GATE_WAIT,
inner: Mutex::new(Inner::default()),
}
}
#[cfg(test)]
fn with_gate_wait(mut self, gate_wait: Duration) -> Self {
self.gate_wait = gate_wait;
self
}
#[must_use]
pub fn path(&self) -> &Path {
&self.path
}
pub fn try_hold(&self) -> LeaseStatus {
let mut inner = self.inner.lock();
if inner.file.is_some() {
return LeaseStatus::Held;
}
match self.acquire() {
Ok(file) => {
inner.file = Some(file);
inner.last_denial = None;
tracing::warn!(
pid = std::process::id(),
lock = %self.path.display(),
"maintenance lease acquired: this process runs dream and TTL-purge \
passes for this data root (#8733)"
);
LeaseStatus::Held
}
Err(status) => {
if inner.last_denial.as_ref() != Some(&status) {
log_denial(&self.path, &status);
inner.last_denial = Some(status.clone());
}
status
}
}
}
fn acquire(&self) -> Result<File, LeaseStatus> {
let unavailable = |e: std::io::Error| LeaseStatus::Unavailable {
reason: e.to_string(),
};
let mut file = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&self.path)
.map_err(unavailable)?;
let _gate = enter_pid_gate(&gate_path(&self.path), self.gate_wait);
match file.try_lock() {
Ok(()) => {}
Err(TryLockError::WouldBlock) => {
return Err(LeaseStatus::HeldElsewhere {
holder_pid: recorded_holder_pid(&self.path),
});
}
Err(TryLockError::Error(e)) => return Err(unavailable(e)),
}
let pid_line = format!("{}\n", std::process::id());
let recorded = file
.set_len(0)
.and_then(|()| file.seek(SeekFrom::Start(0)))
.and_then(|_| file.write_all(pid_line.as_bytes()));
if let Err(e) = recorded {
tracing::warn!(lock = %self.path.display(), "maintenance lease: pid not recorded: {e}");
}
Ok(file)
}
}
fn gate_path(lock_path: &Path) -> PathBuf {
let mut name = lock_path.as_os_str().to_owned();
name.push(PID_GATE_SUFFIX);
PathBuf::from(name)
}
fn enter_pid_gate(path: &Path, wait: Duration) -> Option<File> {
let gate = OpenOptions::new()
.write(true)
.create(true)
.truncate(false)
.open(path)
.inspect_err(|e| tracing::debug!(gate = %path.display(), "pid gate not opened: {e}"))
.ok()?;
let deadline = Instant::now() + wait;
loop {
match gate.try_lock() {
Ok(()) => return Some(gate),
Err(TryLockError::WouldBlock) if Instant::now() < deadline => {
std::thread::sleep(PID_GATE_POLL);
}
Err(e) => {
tracing::debug!(gate = %path.display(), "pid gate not taken: {e}");
return None;
}
}
}
}
fn log_denial(path: &Path, status: &LeaseStatus) {
match status {
LeaseStatus::HeldElsewhere { holder_pid } => tracing::warn!(
holder_pid = ?holder_pid,
lock = %path.display(),
"maintenance lease held by another process: this process serves reads and \
writes but runs no dream or TTL-purge pass (#8733)"
),
LeaseStatus::Unavailable { reason } => tracing::warn!(
lock = %path.display(),
"maintenance lease unavailable ({reason}): this process runs no dream or \
TTL-purge pass (fail closed, #8733)"
),
LeaseStatus::Held => {}
}
}
#[must_use]
pub fn recorded_holder_pid(lock_path: &Path) -> Option<u32> {
std::fs::read_to_string(lock_path).ok()?.trim().parse().ok()
}
#[cfg(test)]
#[path = "maintenance_lease_tests.rs"]
mod maintenance_lease_tests;