use std::time::Duration;
use super::error::StateError;
use super::row::Row;
use super::store::Store;
use crate::process::ProcessStamp;
use crate::types::{Pid, ProcessStartTime};
pub const FOREIGN_STALE_BUDGET: Duration = Duration::from_secs(30 * 60);
#[derive(Debug)]
pub struct LockClaim {
pub instance: String,
pub operation: String,
holder: ProcessStamp,
host: String,
}
impl Store {
pub fn claim_lock(&self, instance: &str, operation: &str) -> Result<LockClaim, StateError> {
let me = ProcessStamp::current();
let host = Self::hostname();
let claimed = self.execute(
"INSERT INTO op_locks
(instance, operation, holder_pid, holder_start_time, holder_host, acquired_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6)
ON CONFLICT(instance) DO UPDATE SET
operation = excluded.operation,
acquired_at = excluded.acquired_at
WHERE op_locks.holder_pid = excluded.holder_pid
AND op_locks.holder_start_time = excluded.holder_start_time
AND op_locks.holder_host = excluded.holder_host",
&[
instance.into(),
operation.into(),
me.pid.into(),
me.start_time.into(),
host.as_str().into(),
Self::now().into(),
],
)?;
if claimed > 0 {
return Ok(LockClaim {
instance: instance.into(),
operation: operation.into(),
holder: me,
host,
});
}
let Some(existing) = self.query_row(
"SELECT operation, holder_pid, holder_start_time, holder_host, acquired_at
FROM op_locks WHERE instance = ?1",
&[instance.into()],
|row: &Row| {
Ok((
row.get_string(0)?,
row.get_u32(1)?,
row.get_i64(2)?,
row.get_string(3)?,
row.get_i64(4)?,
))
},
)?
else {
return self.takeover_or_held(instance, operation, &me, &host, None);
};
let (held_op, holder_pid, holder_start, holder_host, acquired_at) = existing;
let holder = ProcessStamp {
pid: Pid::from_os(holder_pid),
start_time: ProcessStartTime::from_os(holder_start as u64),
};
let same_host = holder_host == host;
let takeover_ok = if same_host {
!holder.is_alive()
} else {
Self::now() - acquired_at >= FOREIGN_STALE_BUDGET.as_secs() as i64
};
if !takeover_ok {
return Err(StateError::LockHeld {
instance: instance.into(),
operation: held_op,
holder_pid,
acquired_at,
});
}
self.takeover_or_held(
instance,
operation,
&me,
&host,
Some((holder_pid, holder_start, holder_host)),
)
}
fn takeover_or_held(
&self,
instance: &str,
operation: &str,
me: &ProcessStamp,
host: &str,
dead_holder: Option<(u32, i64, String)>,
) -> Result<LockClaim, StateError> {
let won = match dead_holder {
Some((dead_pid, dead_start, dead_host)) => self.execute(
"UPDATE op_locks SET
operation = ?2, holder_pid = ?3, holder_start_time = ?4,
holder_host = ?5, acquired_at = ?6
WHERE instance = ?1
AND holder_pid = ?7 AND holder_start_time = ?8 AND holder_host = ?9",
&[
instance.into(),
operation.into(),
me.pid.into(),
me.start_time.into(),
host.into(),
Self::now().into(),
dead_pid.into(),
dead_start.into(),
dead_host.into(),
],
)?,
None => self.execute(
"INSERT INTO op_locks
(instance, operation, holder_pid, holder_start_time, holder_host, acquired_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6)
ON CONFLICT(instance) DO NOTHING",
&[
instance.into(),
operation.into(),
me.pid.into(),
me.start_time.into(),
host.into(),
Self::now().into(),
],
)?,
};
if won > 0 {
Ok(LockClaim {
instance: instance.into(),
operation: operation.into(),
holder: *me,
host: host.to_owned(),
})
} else {
Err(StateError::LockHeld {
instance: instance.into(),
operation: operation.into(),
holder_pid: me.pid.get(),
acquired_at: Self::now(),
})
}
}
pub fn release_lock(&self, claim: &LockClaim) -> Result<(), StateError> {
self.execute(
"DELETE FROM op_locks
WHERE instance = ?1 AND holder_pid = ?2 AND holder_start_time = ?3
AND holder_host = ?4",
&[
claim.instance.as_str().into(),
claim.holder.pid.into(),
claim.holder.start_time.into(),
claim.host.as_str().into(),
],
)?;
Ok(())
}
pub fn lock_holder_alive(&self, instance: &str) -> Result<bool, StateError> {
let host = Self::hostname();
let existing = self.query_row(
"SELECT holder_pid, holder_start_time, holder_host FROM op_locks WHERE instance = ?1",
&[instance.into()],
|row: &Row| Ok((row.get_u32(0)?, row.get_i64(1)?, row.get_string(2)?)),
)?;
Ok(existing.is_some_and(|(pid, start_time, holder_host)| {
let same_host = holder_host == host;
if same_host {
ProcessStamp {
pid: Pid::from_os(pid),
start_time: ProcessStartTime::from_os(start_time as u64),
}
.is_alive()
} else {
true
}
}))
}
}