use std::fmt;
use std::sync::{Arc, Mutex};
use uuid::Uuid;
use crate::config::{DnsConfig, validate_against};
use crate::error::{ConflictReason, Error, Result};
use crate::journal::JournalRecord;
use crate::manager::Inner;
use crate::ownership::{ResourceId, ResourceLock};
use crate::platform::OwnershipProof;
use crate::platform::ResourceStatus;
pub(crate) struct LiveRecord {
pub(crate) record: JournalRecord,
pub(crate) verified: Option<OwnershipProof>,
}
pub(crate) struct LiveLease {
live: Vec<Arc<Mutex<LiveRecord>>>,
_locks: Vec<ResourceLock>,
}
pub struct Lease {
inner: Arc<Inner>,
resources: Vec<ResourceId>,
lease_id: Uuid,
is_noop: bool,
state: Mutex<Option<LiveLease>>,
}
impl Lease {
pub(crate) fn new_owned(
inner: Arc<Inner>,
lease_id: Uuid,
records: Vec<JournalRecord>,
verified: Vec<Option<OwnershipProof>>,
locks: Vec<ResourceLock>,
was_noop: bool,
) -> Self {
let mut resources = Vec::with_capacity(records.len());
let mut live = Vec::with_capacity(records.len());
for (record, verified) in records.into_iter().zip(verified) {
resources.push(record.resource.clone());
let shared = Arc::new(Mutex::new(LiveRecord { record, verified }));
inner.register_active(Arc::clone(&shared));
live.push(shared);
}
Self {
inner,
resources,
lease_id,
is_noop: was_noop,
state: Mutex::new(Some(LiveLease {
live,
_locks: locks,
})),
}
}
pub fn resources(&self) -> &[ResourceId] {
&self.resources
}
pub fn lease_id(&self) -> Uuid {
self.lease_id
}
pub fn is_noop(&self) -> bool {
self.is_noop
}
pub fn update(&self, config: &DnsConfig) -> Result<()> {
let caps = self.inner.backend.capabilities();
let plan = validate_against(config, &caps)?;
self.inner.backend.validate_plan(config.scope(), &plan)?;
let mut guard = self
.state
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let Some(state) = guard.take() else {
return Err(Error::Conflict {
resource: self
.resources
.first()
.cloned()
.ok_or_else(|| Error::invalid_config("lease owns no resources"))?,
reason: ConflictReason::LeaseNotActive,
});
};
let LiveLease { live, _locks } = state;
let repack = |live: Vec<Arc<Mutex<LiveRecord>>>| LiveLease { live, _locks };
for entry in &live {
let record = entry
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.record
.clone();
match self.inner.backend.resource_status(&record.identity) {
Ok(ResourceStatus::Same) => {}
Ok(status @ (ResourceStatus::Gone | ResourceStatus::Replaced)) => {
if let Err(error) = self
.inner
.journal
.remove(&record.lease_id, &record.resource)
{
*guard = Some(repack(live));
return Err(error);
}
self.inner.unregister_active(&record.resource);
*guard = Some(repack(live));
return Err(match status {
ResourceStatus::Gone => Error::ResourceGone {
backend: self.inner.backend.kind(),
resource: record.resource,
message: "the leased native resource incarnation is gone; its journal was cleared and a fresh lease is required".to_string(),
},
ResourceStatus::Replaced => Error::ResourceIdentity {
backend: self.inner.backend.kind(),
resource: record.resource,
message: "the lease target was replaced; its journal was cleared and a fresh lease is required".to_string(),
},
_ => unreachable!(),
});
}
Ok(ResourceStatus::Ambiguous) => {
*guard = Some(repack(live));
return Err(Error::ResourceIdentity {
backend: self.inner.backend.kind(),
resource: record.resource,
message: "the leased native resource incarnation cannot be proven; refusing update".to_string(),
});
}
Err(error) => {
*guard = Some(repack(live));
return Err(error);
}
}
}
let wanted = match self.inner.backend.resolve_resources(config.scope(), &plan) {
Ok(wanted) => wanted,
Err(error) => {
*guard = Some(repack(live));
return Err(error);
}
};
let mut wanted_sorted = wanted;
wanted_sorted.sort();
let mut owned_sorted: Vec<ResourceId> = live
.iter()
.map(|record| {
record
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.record
.resource
.clone()
})
.collect();
owned_sorted.sort();
if wanted_sorted != owned_sorted {
*guard = Some(repack(live));
return Err(Error::UpdateRequiresRebind {
owned: owned_sorted,
requested: wanted_sorted,
});
}
let tokens: Vec<std::sync::Arc<std::sync::Mutex<()>>> = owned_sorted
.iter()
.map(|resource| self.inner.lease_token(resource))
.collect();
let token_guards: Vec<_> = tokens
.iter()
.map(|token| {
token
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
})
.collect();
let result = self.inner.transact_update(&live, &plan);
drop(token_guards);
drop(tokens);
*guard = Some(repack(live));
result
}
#[allow(clippy::result_large_err)]
pub fn restore(self) -> std::result::Result<(), RestoreFailure> {
match self.restore_state() {
Ok(()) => Ok(()),
Err(error) => Err(RestoreFailure { error, lease: self }),
}
}
fn restore_state(&self) -> Result<()> {
let mut guard = self
.state
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let Some(state) = guard.take() else {
return Err(Error::Conflict {
resource: self
.resources
.first()
.cloned()
.ok_or_else(|| Error::invalid_config("lease owns no resources"))?,
reason: ConflictReason::LeaseNotActive,
});
};
let LiveLease { live, _locks } = state;
let mut first_error = None;
for record in &live {
self.inner.with_live_record(record, |live| {
let resource = live.record.resource.clone();
if let Err(error) = self.inner.finalize_live(live, None) {
if first_error.is_none() {
first_error = Some(error);
}
return;
}
match self.inner.restore_lease_state(&live.record) {
Ok(()) => self.inner.unregister_active(&resource),
Err(error) => {
if first_error.is_none() {
first_error = Some(error);
}
}
}
});
}
match first_error {
None => {
drop(_locks);
self.inner.release_enforce_watch();
Ok(())
}
Some(error) => {
*guard = Some(LiveLease { live, _locks });
Err(error)
}
}
}
pub fn abandon(self) -> Result<()> {
let mut guard = self
.state
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if let Some(LiveLease { live, _locks }) = guard.take() {
let mut failure = None;
for record in &live {
self.inner.with_live_record(record, |live| {
let resource = live.record.resource.clone();
if failure.is_none()
&& let Err(error) =
self.inner.journal.remove(&live.record.lease_id, &resource)
{
failure = Some(error);
}
self.inner.unregister_active(&resource);
});
}
drop(_locks);
self.inner.release_enforce_watch();
if let Some(error) = failure {
return Err(error);
}
}
Ok(())
}
#[cfg(feature = "test-util")]
pub fn debug_release_locks_keep_journal(self) {
let mut guard = self
.state
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if let Some(LiveLease { live, _locks }) = guard.take() {
for record in &live {
self.inner.with_live_record(record, |live| {
let resource = live.record.resource.clone();
self.inner.unregister_active(&resource);
});
}
drop(_locks);
self.inner.release_enforce_watch();
}
}
}
pub struct RestoreFailure {
pub error: Error,
pub lease: Lease,
}
impl From<RestoreFailure> for Error {
fn from(failure: RestoreFailure) -> Self {
failure.error
}
}
impl fmt::Debug for RestoreFailure {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("RestoreFailure")
.field("error", &self.error)
.field("lease", &self.lease)
.finish()
}
}
impl fmt::Debug for Lease {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Lease")
.field("owner", &self.inner.owner)
.field("resources", &self.resources)
.field("lease_id", &self.lease_id)
.field("is_noop", &self.is_noop)
.finish()
}
}
impl Drop for Lease {
fn drop(&mut self) {
if let Some(LiveLease { live, _locks }) = self
.state
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.take()
{
for record in &live {
self.inner.with_live_record(record, |live| {
let resource = live.record.resource.clone();
let _ = self.inner.finalize_live(live, None);
self.inner.best_effort_restore(&live.record);
self.inner.unregister_active(&resource);
});
}
drop(_locks);
self.inner.release_enforce_watch();
}
}
}