use std::sync::atomic::Ordering;
use std::sync::{Arc, Weak};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use nodedb_cluster::{METADATA_GROUP_ID, MetadataEntry, WaitOutcome, encode_entry};
#[cfg(test)]
use nodedb_cluster::AppliedIndexWatcher;
use crate::control::catalog_entry::{self, CatalogEntry};
use crate::control::state::SharedState;
use crate::error::Error;
pub const DEFAULT_PROPOSE_TIMEOUT: Duration = Duration::from_secs(5);
pub const DEFAULT_DRAIN_TIMEOUT: Duration = Duration::from_secs(35);
const DDL_PREPARE_LEASE: Duration = Duration::from_secs(60);
const DDL_PREPARE_WAIT: Duration = Duration::from_secs(70);
pub trait MetadataRaftHandle: Send + Sync {
fn propose(&self, bytes: Vec<u8>) -> Result<u64, Error>;
}
pub struct RaftLoopProposerHandle {
raft_loop: Weak<
nodedb_cluster::RaftLoop<
crate::control::cluster::SpscCommitApplier,
crate::control::LocalPlanExecutor,
>,
>,
}
impl RaftLoopProposerHandle {
pub fn new(
raft_loop: Arc<
nodedb_cluster::RaftLoop<
crate::control::cluster::SpscCommitApplier,
crate::control::LocalPlanExecutor,
>,
>,
) -> Self {
Self {
raft_loop: Arc::downgrade(&raft_loop),
}
}
}
impl MetadataRaftHandle for RaftLoopProposerHandle {
fn propose(&self, bytes: Vec<u8>) -> Result<u64, Error> {
let raft_loop = self.raft_loop.upgrade().ok_or_else(|| Error::Config {
detail: "metadata propose: cluster not running".into(),
})?;
tokio::task::block_in_place(|| {
tokio::runtime::Handle::current()
.block_on(raft_loop.propose_to_metadata_group_via_leader(bytes))
})
.map_err(|e| Error::Config {
detail: format!("metadata propose: {e}"),
})
}
}
fn wall_now_ns() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_nanos()
.min(u64::MAX as u128) as u64
}
fn propose_metadata_and_wait(
shared: &SharedState,
handle: &dyn MetadataRaftHandle,
entry: &MetadataEntry,
timeout: Duration,
) -> Result<u64, Error> {
let raw = encode_entry(entry).map_err(|e| Error::Config {
detail: format!("metadata entry encode: {e}"),
})?;
let index = handle.propose(raw)?;
let watcher = shared.applied_index_watcher(METADATA_GROUP_ID);
let outcome = tokio::task::block_in_place(|| watcher.wait_for(index, timeout));
match outcome {
WaitOutcome::Reached => Ok(index),
WaitOutcome::TimedOut => Err(Error::Config {
detail: format!(
"metadata propose timed out after {timeout:?} waiting for log index {index} (current: {})",
watcher.current()
),
}),
WaitOutcome::GroupGone => Err(Error::Config {
detail: "metadata group no longer hosted on this node".into(),
}),
}
}
pub(crate) struct DdlPrepareGuard<'a> {
shared: &'a SharedState,
handle: &'a dyn MetadataRaftHandle,
token: u64,
}
impl DdlPrepareGuard<'_> {
pub(crate) fn token(&self) -> u64 {
self.token
}
}
impl Drop for DdlPrepareGuard<'_> {
fn drop(&mut self) {
if let Err(error) = propose_metadata_and_wait(
self.shared,
self.handle,
&MetadataEntry::DdlPrepareRelease { token: self.token },
DEFAULT_PROPOSE_TIMEOUT,
) {
tracing::error!(token = self.token, %error, "metadata DDL lease release failed");
}
}
}
pub(crate) fn acquire_ddl_prepare_lease<'a>(
shared: &'a SharedState,
handle: &'a dyn MetadataRaftHandle,
) -> Result<DdlPrepareGuard<'a>, Error> {
let sequence = shared
.metadata_ddl_token_seq
.fetch_add(1, Ordering::Relaxed);
let token = shared.node_id.wrapping_mul(0x9e37_79b9_7f4a_7c15)
^ wall_now_ns().rotate_left(17)
^ sequence;
let deadline = Instant::now() + DDL_PREPARE_WAIT;
loop {
propose_metadata_and_wait(
shared,
handle,
&MetadataEntry::DdlPrepareAcquire { token },
DEFAULT_PROPOSE_TIMEOUT,
)?;
loop {
let owner = *shared
.metadata_ddl_owner
.lock()
.map_err(|_| Error::Config {
detail: "metadata DDL owner lock poisoned".into(),
})?;
match owner {
Some((current, _)) if current == token => {
return Ok(DdlPrepareGuard {
shared,
handle,
token,
});
}
Some((current, acquired_at))
if shared.is_metadata_leader()
&& acquired_at.elapsed() >= DDL_PREPARE_LEASE =>
{
propose_metadata_and_wait(
shared,
handle,
&MetadataEntry::DdlPrepareRelease { token: current },
DEFAULT_PROPOSE_TIMEOUT,
)?;
break;
}
None => break,
Some(_) if Instant::now() < deadline => {
std::thread::sleep(Duration::from_millis(10));
}
Some(_) => {
return Err(Error::Config {
detail: "metadata DDL preparation lease timed out".into(),
});
}
}
}
}
}
pub fn propose_catalog_entry(shared: &SharedState, entry: &CatalogEntry) -> Result<u64, Error> {
propose_catalog_entry_with_timeout(shared, entry, DEFAULT_PROPOSE_TIMEOUT)
}
pub fn propose_catalog_entry_with_timeout(
shared: &SharedState,
entry: &CatalogEntry,
timeout: Duration,
) -> Result<u64, Error> {
let Some(handle) = shared.metadata_raft.get() else {
return Ok(0);
};
{
let vs = shared.cluster_version_view();
if !vs.can_activate_feature(crate::control::rolling_upgrade::DISTRIBUTED_CATALOG_VERSION) {
tracing::warn!(
min_version = vs.min_version,
required = crate::control::rolling_upgrade::DISTRIBUTED_CATALOG_VERSION,
"metadata propose: cluster in compat mode (mixed-version), \
falling back to legacy direct-write path"
);
return Ok(0);
}
}
if crate::control::server::shared::session::ddl_buffer::try_buffer(entry.clone()) {
return Ok(0);
}
let _local_ddl_guard = shared.metadata_ddl_lock.lock().map_err(|_| Error::Config {
detail: "metadata DDL preparation lock poisoned".into(),
})?;
let distributed_ddl_guard = acquire_ddl_prepare_lease(shared, handle.as_ref())?;
if let Some((descriptor_id, prior_version)) =
crate::control::lease::descriptor_id_and_prior_version(entry, shared)
&& prior_version > 0
{
crate::control::lease::drain_for_ddl(
shared,
descriptor_id,
prior_version,
DEFAULT_DRAIN_TIMEOUT,
)?;
}
let stamped_owned;
let entry: &CatalogEntry = if shared
.cluster_version_view()
.can_activate_feature(crate::control::rolling_upgrade::DESCRIPTOR_VERSIONING_VERSION)
{
stamped_owned = catalog_entry::descriptor_stamp::stamp(
entry.clone(),
&shared.hlc_clock,
shared.credentials.catalog(),
);
&stamped_owned
} else {
entry
};
let payload = catalog_entry::encode(entry)?;
let catalog_entry = match crate::control::server::shared::session::audit_context::current() {
Some(ctx) => MetadataEntry::CatalogDdlAudited {
payload,
auth_user_id: ctx.auth_user_id,
auth_user_name: ctx.auth_user_name,
sql_text: ctx.sql_text,
},
None => MetadataEntry::CatalogDdl { payload },
};
let metadata_entry = MetadataEntry::DdlPrepared {
token: distributed_ddl_guard.token(),
entry: Box::new(catalog_entry),
};
let raw = encode_entry(&metadata_entry).map_err(|e| Error::Config {
detail: format!("metadata entry encode: {e}"),
})?;
let log_index = handle.propose(raw)?;
let watcher = shared.applied_index_watcher(METADATA_GROUP_ID);
let outcome = tokio::task::block_in_place(|| watcher.wait_for(log_index, timeout));
match outcome {
WaitOutcome::Reached
if shared.metadata_ddl_applied_token.load(Ordering::Acquire)
== distributed_ddl_guard.token() =>
{
Ok(log_index)
}
WaitOutcome::Reached => Err(Error::Config {
detail: "metadata DDL preparation ownership was superseded before apply".into(),
}),
WaitOutcome::TimedOut => Err(Error::Config {
detail: format!(
"metadata propose timed out after {:?} waiting for log index {} (current: {})",
timeout,
log_index,
watcher.current()
),
}),
WaitOutcome::GroupGone => Err(Error::Config {
detail: "metadata group no longer hosted on this node".into(),
}),
}
}
pub fn propose_surrogate_hwm(shared: &SharedState, hwm: u32) -> Result<u64, Error> {
let Some(handle) = shared.metadata_raft.get() else {
return Ok(0);
};
let entry = MetadataEntry::SurrogateAlloc { hwm };
let raw = encode_entry(&entry).map_err(|e| Error::Config {
detail: format!("surrogate_alloc encode: {e}"),
})?;
let log_index = handle.propose(raw)?;
let watcher = shared.applied_index_watcher(METADATA_GROUP_ID);
let outcome =
tokio::task::block_in_place(|| watcher.wait_for(log_index, DEFAULT_PROPOSE_TIMEOUT));
if !outcome.is_reached() {
return Err(Error::Config {
detail: format!("surrogate_alloc propose timed out waiting for log index {log_index}"),
});
}
Ok(log_index)
}
pub fn propose_surrogate_reserve(
shared: &SharedState,
node_id: u64,
request_id: u64,
batch_size: u32,
) -> Result<u64, Error> {
let Some(handle) = shared.metadata_raft.get() else {
return Ok(0);
};
let entry = MetadataEntry::SurrogateReserve {
node_id,
request_id,
batch_size,
};
let raw = encode_entry(&entry).map_err(|e| Error::Config {
detail: format!("surrogate_reserve encode: {e}"),
})?;
let log_index = handle.propose(raw)?;
let watcher = shared.applied_index_watcher(METADATA_GROUP_ID);
let outcome =
tokio::task::block_in_place(|| watcher.wait_for(log_index, DEFAULT_PROPOSE_TIMEOUT));
if !outcome.is_reached() {
return Err(Error::Config {
detail: format!(
"surrogate_reserve propose timed out waiting for log index {log_index}"
),
});
}
Ok(log_index)
}
pub fn propose_sync_producer_register(
shared: &SharedState,
lite_id: &str,
producer_id: u64,
tenant_id: u64,
epoch: u64,
created_ms: i64,
) -> Result<u64, Error> {
let Some(handle) = shared.metadata_raft.get() else {
return Ok(0);
};
let entry = MetadataEntry::SyncProducerRegister {
lite_id: lite_id.to_owned(),
producer_id,
tenant_id,
epoch,
created_ms,
};
let raw = encode_entry(&entry).map_err(|e| Error::Config {
detail: format!("sync_producer_register encode: {e}"),
})?;
let log_index = handle.propose(raw)?;
let watcher = shared.applied_index_watcher(METADATA_GROUP_ID);
let outcome =
tokio::task::block_in_place(|| watcher.wait_for(log_index, DEFAULT_PROPOSE_TIMEOUT));
if !outcome.is_reached() {
return Err(Error::Config {
detail: format!(
"sync_producer_register propose timed out waiting for log index {log_index}"
),
});
}
Ok(log_index)
}
pub fn propose_sync_producer_fence(
shared: &SharedState,
lite_id: &str,
new_epoch: u64,
) -> Result<u64, Error> {
let Some(handle) = shared.metadata_raft.get() else {
return Ok(0);
};
let entry = MetadataEntry::SyncProducerFence {
lite_id: lite_id.to_owned(),
new_epoch,
};
let raw = encode_entry(&entry).map_err(|e| Error::Config {
detail: format!("sync_producer_fence encode: {e}"),
})?;
let log_index = handle.propose(raw)?;
let watcher = shared.applied_index_watcher(METADATA_GROUP_ID);
let outcome =
tokio::task::block_in_place(|| watcher.wait_for(log_index, DEFAULT_PROPOSE_TIMEOUT));
if !outcome.is_reached() {
return Err(Error::Config {
detail: format!(
"sync_producer_fence propose timed out waiting for log index {log_index}"
),
});
}
Ok(log_index)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn watcher_helper_returns_reached_on_past_target() {
let w = AppliedIndexWatcher::new();
w.bump(10);
assert!(w.wait_for(5, Duration::from_millis(1)).is_reached());
}
}