use crate::memory_core::maintenance_log::{self, DeletionReason, MaintenanceDeletion};
use crate::memory_core::palace::PalaceId;
use crate::memory_core::store::kg_redb::KgStoreRedb;
use std::sync::Arc;
use std::thread::Builder as OpenWriteThread;
use std::time::Duration;
use uuid::Uuid;
pub(crate) const OPEN_WRITE_BUDGET: Duration = Duration::from_secs(2);
pub(crate) fn run_open_write<R, F>(
store: &KgStoreRedb,
palace: &PalaceId,
operation: &'static str,
budget: Duration,
work: F,
) -> Option<R>
where
R: Send + 'static,
F: FnOnce() -> R + Send + 'static,
{
let Some(claim) = store.try_claim_open_write() else {
tracing::warn!(
palace = %palace,
operation,
"#8314: skipped: an earlier open-time kg.redb write for this palace \
has not finished; it re-runs on a later open"
);
return None;
};
let (done_tx, done_rx) = std::sync::mpsc::sync_channel::<R>(1);
let spawned = OpenWriteThread::new()
.name("trusty-open-write".into())
.spawn(move || {
let out = work();
drop(claim);
let _ = done_tx.send(out);
});
if let Err(e) = spawned {
tracing::warn!(palace = %palace, operation, "open-time write thread failed to start: {e}");
return None;
}
match done_rx.recv_timeout(budget) {
Ok(out) => Some(out),
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
tracing::error!(
palace = %palace,
operation,
budget_ms = budget.as_millis(),
"#8314: open-time write waited past its budget on a kg.redb \
write that has not finished; the open proceeds and the write \
lands when that write ends"
);
None
}
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
tracing::warn!(palace = %palace, operation, "open-time write panicked");
None
}
}
}
pub(super) fn reclaim_expired_rows(
store: Arc<KgStoreRedb>,
ids: Vec<Uuid>,
palace: &PalaceId,
data_dir: Option<std::path::PathBuf>,
budget: Duration,
) -> usize {
if ids.is_empty() {
return 0;
}
let owner = palace.clone();
let worker = Arc::clone(&store);
let work = move || {
let mut reclaimed = 0usize;
for id in ids {
match worker.delete_drawer(id) {
Ok(()) => {
reclaimed += 1;
let rec =
MaintenanceDeletion::new(&owner, id, DeletionReason::ExpiredPurgeAtOpen);
maintenance_log::record(data_dir.as_deref(), &rec);
}
Err(e) => tracing::warn!(
palace = %owner, id = %id,
"purge_expired: delete_drawer failed: {e:#}"
),
}
}
reclaimed
};
run_open_write(&store, palace, "purge_expired_at_open", budget, work).unwrap_or(0)
}