use std::sync::{Arc, Mutex};
use std::thread::JoinHandle;
use std::time::{Duration, Instant};
use dashmap::DashMap;
use rust_rocksdb::WriteBatch;
use super::engine::META_CF;
use super::RocksDb as DB;
use crate::error::{DbError, DbResult};
const MARKER_PREFIX: &str = "pending_drop:";
#[derive(Clone, Copy, PartialEq)]
enum DropState {
Pending,
Dropping,
}
pub enum Claim {
Claimed,
InProgress,
NotPending,
}
#[derive(Default)]
pub struct PendingCfDrops {
states: DashMap<String, DropState>,
droppers: Mutex<Vec<JoinHandle<()>>>,
}
impl PendingCfDrops {
pub fn new() -> Arc<Self> {
Arc::new(Self::default())
}
pub fn contains(&self, cf_name: &str) -> bool {
self.states.contains_key(cf_name)
}
pub fn schedule(&self, db: &DB, db_meta_key: &str, cfs: &[String]) -> DbResult<()> {
let meta_cf = db
.cf_handle(META_CF)
.ok_or_else(|| DbError::InternalError("_meta column family missing".to_string()))?;
let mut batch = WriteBatch::default();
batch.delete_cf(&meta_cf, db_meta_key.as_bytes());
for cf in cfs {
batch.put_cf(
&meta_cf,
format!("{}{}", MARKER_PREFIX, cf).as_bytes(),
b"1",
);
}
db.write(&batch).map_err(|e| {
DbError::InternalError(format!("Failed to schedule collection drops: {}", e))
})?;
for cf in cfs {
self.states.insert(cf.clone(), DropState::Pending);
}
Ok(())
}
pub fn resume_from_meta(&self, db: &DB) -> Vec<String> {
let meta_cf = match db.cf_handle(META_CF) {
Some(cf) => cf,
None => return vec![],
};
let iter = db.prefix_iterator_cf(&meta_cf, MARKER_PREFIX.as_bytes());
let cfs: Vec<String> = iter
.filter_map(|result| {
result.ok().and_then(|(key, _)| {
let key_str = String::from_utf8(key.to_vec()).ok()?;
key_str.strip_prefix(MARKER_PREFIX).map(|s| s.to_string())
})
})
.collect();
for cf in &cfs {
self.states.insert(cf.clone(), DropState::Pending);
}
cfs
}
pub fn claim_for_recreate(&self, cf_name: &str) -> Claim {
use dashmap::mapref::entry::Entry;
match self.states.entry(cf_name.to_string()) {
Entry::Occupied(mut entry) => match entry.get() {
DropState::Pending => {
*entry.get_mut() = DropState::Dropping;
Claim::Claimed
}
DropState::Dropping => Claim::InProgress,
},
Entry::Vacant(_) => Claim::NotPending,
}
}
pub fn release_claim(&self, cf_name: &str) {
self.states.insert(cf_name.to_string(), DropState::Pending);
}
pub fn complete(&self, db: &DB, cf_name: &str) {
if let Some(meta_cf) = db.cf_handle(META_CF) {
let _ = db.delete_cf(&meta_cf, format!("{}{}", MARKER_PREFIX, cf_name).as_bytes());
}
self.states.remove(cf_name);
super::collection::index_meta::invalidate_index_meta(db, cf_name);
}
pub fn wait_until_dropped(&self, cf_name: &str, timeout: Duration) -> DbResult<()> {
let start = Instant::now();
while self.states.contains_key(cf_name) {
if start.elapsed() > timeout {
return Err(DbError::InternalError(format!(
"Timed out waiting for pending drop of collection '{}'",
cf_name
)));
}
std::thread::sleep(Duration::from_millis(10));
}
Ok(())
}
fn begin_drop(&self, cf_name: &str) -> bool {
use dashmap::mapref::entry::Entry;
match self.states.entry(cf_name.to_string()) {
Entry::Occupied(mut entry) if *entry.get() == DropState::Pending => {
*entry.get_mut() = DropState::Dropping;
true
}
_ => false,
}
}
pub fn spawn_dropper(db: Arc<DB>, registry: Arc<Self>, cfs: Vec<String>) {
if cfs.is_empty() {
return;
}
let registry_for_handle = Arc::clone(®istry);
let handle = std::thread::spawn(move || {
let start = Instant::now();
let total = cfs.len();
let mut dropped = 0usize;
for (i, cf) in cfs.iter().enumerate() {
if i > 0 {
std::thread::sleep(Duration::from_millis(25));
}
if !registry.begin_drop(cf) {
continue; }
if db.cf_handle(cf).is_some() {
if let Err(e) = super::cf_ops::timed(|| db.drop_cf(cf)) {
tracing::warn!("Background drop of column family '{}' failed: {}", cf, e);
registry.release_claim(cf);
continue;
}
dropped += 1;
}
registry.complete(&db, cf);
}
tracing::info!(
"Background-dropped {}/{} column families in {:.2?}",
dropped,
total,
start.elapsed()
);
});
let locked = registry_for_handle.droppers.lock();
if let Ok(mut handles) = locked {
handles.retain(|h| !h.is_finished());
handles.push(handle);
}
}
pub fn join_droppers(&self) {
let handles = match self.droppers.lock() {
Ok(mut guard) => std::mem::take(&mut *guard),
Err(_) => return,
};
for handle in handles {
let _ = handle.join();
}
}
}