use std::collections::BTreeSet;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use tokio::sync::mpsc::UnboundedSender;
use crate::bridge::envelope::PhysicalPlan;
use crate::control::cluster::calvin::scheduler::driver::core::routing::{PlanRouting, plan_vshard};
use crate::control::cluster::calvin::scheduler::lock_manager::{LockKey, LockManager, TxnId};
use crate::control::planner::calvin::is_dependent_predicate;
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TenantId, VShardId};
use nodedb_physical::physical_plan::MetaOp;
use super::lock_keys::plan_lock_keys;
use super::predicate::plan_is_write;
use super::write_order_lock::KeyedWriteOrderLock;
static ROUTED_TO_CALVIN: AtomicU64 = AtomicU64::new(0);
pub fn cp_routed_to_calvin() -> u64 {
ROUTED_TO_CALVIN.load(Ordering::Relaxed)
}
pub struct WriteTarget<'a> {
pub tenant_id: TenantId,
pub database_id: DatabaseId,
pub vshard_id: VShardId,
pub plan: &'a PhysicalPlan,
}
pub enum WriteAdmission {
FastPath { guard: Option<WriteAdmissionGuard> },
FastPathBlocking {
key: LockKey,
keyed_lock: Arc<KeyedWriteOrderLock>,
},
RouteToCalvin,
ExemptRead,
}
pub struct WriteAdmissionGuard {
lock_manager: Arc<Mutex<LockManager>>,
txn: TxnId,
promotion_sender: Option<UnboundedSender<Vec<TxnId>>>,
}
impl Drop for WriteAdmissionGuard {
fn drop(&mut self) {
let promoted = self
.lock_manager
.lock()
.unwrap_or_else(|p| p.into_inner())
.release(self.txn);
if !promoted.is_empty()
&& let Some(sender) = &self.promotion_sender
&& let Err(e) = sender.send(promoted)
{
tracing::warn!(
error = %e,
"write-admission gate: could not deliver promoted Calvin waiters to \
the scheduler (receiver gone); those transactions may stall"
);
}
}
}
pub fn admit(shared: &SharedState, target: &WriteTarget<'_>) -> WriteAdmission {
if matches!(
target.plan,
PhysicalPlan::Meta(
MetaOp::CalvinExecuteStatic { .. }
| MetaOp::CalvinExecuteActive { .. }
| MetaOp::CalvinFlush { .. }
| MetaOp::CalvinDrop { .. }
| MetaOp::CalvinResolve { .. }
)
) {
return WriteAdmission::ExemptRead;
}
if !plan_is_write(target.plan) {
return WriteAdmission::ExemptRead;
}
let point_keys = plan_lock_keys(target.plan);
let is_predicate = is_dependent_predicate(target.plan);
let vshard = match &point_keys {
Some((v, _)) => *v,
None if is_predicate => match plan_vshard(target.plan) {
PlanRouting::Vshards(v) => match v.as_slice() {
[v] => *v,
_ => return WriteAdmission::FastPath { guard: None },
},
PlanRouting::ControlPlaneOnly | PlanRouting::NotAWrite | PlanRouting::Unroutable(_) => {
return WriteAdmission::FastPath { guard: None };
}
},
None => return WriteAdmission::FastPath { guard: None },
};
let Some(lock_manager) = shared
.calvin_lock_managers
.lock()
.unwrap_or_else(|p| p.into_inner())
.get(&vshard.as_u32())
.map(Arc::clone)
else {
return match point_keys.and_then(|(_v, keys)| single_point_key(keys)) {
Some(key) => WriteAdmission::FastPathBlocking {
key,
keyed_lock: Arc::clone(&shared.write_order_locks),
},
None => WriteAdmission::FastPath { guard: None },
};
};
let Some((_v, keys)) = point_keys else {
ROUTED_TO_CALVIN.fetch_add(1, Ordering::Relaxed);
return WriteAdmission::RouteToCalvin;
};
let txn = TxnId::new(
TxnId::AUTOCOMMIT_EPOCH,
shared.autocommit_lock_seq.fetch_add(1, Ordering::Relaxed),
);
let acquired = {
let mut lm = lock_manager.lock().unwrap_or_else(|p| p.into_inner());
lm.try_acquire(txn, keys)
};
if acquired {
let promotion_sender = shared
.calvin_promotion_senders
.lock()
.unwrap_or_else(|p| p.into_inner())
.get(&vshard.as_u32())
.cloned();
WriteAdmission::FastPath {
guard: Some(WriteAdmissionGuard {
lock_manager,
txn,
promotion_sender,
}),
}
} else {
ROUTED_TO_CALVIN.fetch_add(1, Ordering::Relaxed);
WriteAdmission::RouteToCalvin
}
}
fn single_point_key(keys: BTreeSet<LockKey>) -> Option<LockKey> {
let mut it = keys.into_iter();
match (it.next(), it.next()) {
(Some(key), None) => Some(key),
_ => None,
}
}