use std::collections::BTreeSet;
use super::*;
const STATUS_SCHEMA_VERSION: u32 = 5;
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct HostLeasePendingRequest {
pub waiter_id: String,
pub priority_class: HostLeasePriorityClass,
pub requested_at_ms: i64,
pub deadline_at_ms: i64,
pub owner_pid: Option<u32>,
pub recoverable: bool,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct HostLeaseState {
pub schema_version: u32,
pub host: String,
#[serde(default)]
pub resource_class: HostLeaseResourceClass,
#[serde(default = "default_host_lease_domain")]
pub domain: String,
pub observed_at_ms: i64,
#[serde(default)]
pub active: Option<HostLeaseHandle>,
pub pending: Vec<HostLeasePendingRequest>,
pub recovered_stale_lease: bool,
#[serde(default)]
pub recovered: Option<HostLeaseHandle>,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct HostLeaseOverview {
pub schema_version: u32,
pub host: String,
pub observed_at_ms: i64,
pub resources: Vec<HostLeaseState>,
}
impl HostLeaseStore {
pub fn status_overview(
&self,
host: &str,
domain: Option<&str>,
) -> Result<HostLeaseOverview, HostLeaseError> {
let host = normalize_component("host", host)?;
let domain = domain.map(normalize_domain).transpose()?;
let mut conn = self.connection(SQLITE_MUTATION_BUSY_TIMEOUT)?;
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let now = unix_now_ms()?;
let removed = admission::cleanup_waiters(&tx, now, self.process_inspector.as_ref())?;
let keys = resource_keys(&tx, &host, domain.as_deref())?;
let resources = keys
.into_iter()
.map(|resource| self.observe_resource(&tx, &resource, now))
.collect::<Result<Vec<_>, _>>()?;
let recovered = resources.iter().any(|state| state.recovered.is_some());
tx.commit()?;
if removed || recovered {
self.signal_waiters();
}
Ok(HostLeaseOverview {
schema_version: STATUS_SCHEMA_VERSION,
host,
observed_at_ms: now,
resources,
})
}
pub(super) fn status_in_transaction(
&self,
tx: Transaction<'_>,
host: &str,
resource_class: HostLeaseResourceClass,
domain: &str,
now: i64,
) -> Result<HostLeaseState, HostLeaseError> {
let removed = admission::cleanup_waiters(&tx, now, self.process_inspector.as_ref())?;
let resource = HostLeaseResourceKey {
machine: host.to_string(),
resource_class,
domain: domain.to_string(),
};
let state = self.observe_resource(&tx, &resource, now)?;
tx.commit()?;
if removed || state.recovered.is_some() {
self.signal_waiters();
}
Ok(state)
}
fn observe_resource(
&self,
tx: &Transaction<'_>,
resource: &HostLeaseResourceKey,
now: i64,
) -> Result<HostLeaseState, HostLeaseError> {
let (active, recovered) = active_handle(
tx,
&resource.machine,
resource.resource_class,
&resource.domain,
now,
self.process_inspector.as_ref(),
)?;
Ok(HostLeaseState {
schema_version: STATUS_SCHEMA_VERSION,
host: resource.machine.clone(),
resource_class: resource.resource_class,
domain: resource.domain.clone(),
observed_at_ms: now,
active,
pending: admission::pending_requests(tx, resource)?,
recovered_stale_lease: recovered.is_some(),
recovered,
})
}
}
fn resource_keys(
tx: &Transaction<'_>,
host: &str,
domain: Option<&str>,
) -> Result<Vec<HostLeaseResourceKey>, HostLeaseError> {
let default_domain = domain.unwrap_or(DEFAULT_HOST_LEASE_DOMAIN);
let mut keys: BTreeSet<(String, String)> = HOST_LEASE_RESOURCE_DEFINITIONS
.iter()
.map(|definition| (definition.name.to_string(), default_domain.to_string()))
.collect();
let mut statement = tx.prepare(
"SELECT resource_class, domain FROM host_leases WHERE host = ?1
UNION SELECT resource_class, domain FROM host_lease_waiters WHERE host = ?1",
)?;
for row in statement.query_map(params![host], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
})? {
let (class, stored_domain) = row?;
if domain.is_none_or(|domain| domain == stored_domain) {
keys.insert((class, stored_domain));
}
}
keys.into_iter()
.map(|(class, domain)| {
Ok(HostLeaseResourceKey {
machine: host.to_string(),
resource_class: HostLeaseResourceClass::parse(&class)?,
domain,
})
})
.collect()
}