use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use mongreldb_consensus::engine_sink::{
open_engine_sink, EngineApplySink, EngineGroupConfig, EngineSinkError, EngineTabletPin,
TabletDataCommand, TabletDataCommandRecord, TabletDataMutation, COMMAND_TYPE_TABLET_DATA,
};
use mongreldb_consensus::error::ConsensusError;
use mongreldb_consensus::group::{ConsensusGroup, GroupCommitReceipt, GroupConfig, GroupMetrics};
use mongreldb_consensus::identity::{raft_node_id, CommandKind, RaftNodeId};
use mongreldb_consensus::raft_log::RaftCommitLog;
use mongreldb_consensus::state_machine::ApplySink;
use mongreldb_log::commit_log::{CommitLog, ExecutionControl, LogPosition};
use mongreldb_log::envelope::CommandEnvelope;
use mongreldb_types::hlc::HlcTimestamp;
use mongreldb_types::ids::{DatabaseId, MetadataVersion, NodeId, RaftGroupId, TabletId};
use openraft::BasicNode;
#[cfg(test)]
use serde::{Deserialize, Serialize};
use crate::merge::{
merge_progress, MergeExecutor, MergeInputs, MergeMetaPlane, MergePhase, MergePlan,
MergePlanner, MergePublishCommand,
};
use crate::meta::{MetaCommand, MetaError, MetaGroup, MetaGroupConfig, MetaRejectionReason};
use crate::network::{
InternalRpcHandler, PeerEndpoint, TcpTransport, TransportConfig, TransportError,
TransportSecurity, TransportServer,
};
use crate::node::{ClusterError, NodeIdentity};
use crate::split::{
abort_split, split_progress, ChildAllocation, ChildStateSink, SnapshotPin, SplitAbortReport,
SplitError, SplitExecutor, SplitKeySelection, SplitPhase, SplitPlan, SplitPublishCommand,
TabletDataError, TabletKeyspace, TabletMetaPlane, TabletMutation, TabletSplitPlanner,
};
use crate::tablet::{
Key, ReplicaDescriptor, TablePartitioningRecord, TabletDescriptor, TabletError, TabletLayout,
TabletOwnershipGuard, TabletOwnershipRegistry, TabletState, TABLETS_DIR, TABLET_META_FILENAME,
};
#[derive(Debug, thiserror::Error)]
pub enum RuntimeError {
#[error(transparent)]
Cluster(#[from] ClusterError),
#[error(transparent)]
Meta(#[from] MetaError),
#[error(transparent)]
Consensus(#[from] ConsensusError),
#[error(transparent)]
Transport(#[from] TransportError),
#[error(transparent)]
Tablet(#[from] TabletError),
#[error(transparent)]
EngineSink(#[from] EngineSinkError),
#[error(transparent)]
Split(#[from] SplitError),
#[error(transparent)]
Merge(#[from] crate::merge::MergeError),
#[error("runtime I/O error: {0}")]
Io(#[from] std::io::Error),
#[error("invalid runtime request: {0}")]
InvalidRequest(String),
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct GroupTiming {
pub heartbeat_interval: Duration,
pub election_timeout_min: Duration,
pub election_timeout_max: Duration,
pub install_snapshot_timeout: Duration,
}
impl GroupTiming {
fn apply(&self, config: &mut GroupConfig) {
config.heartbeat_interval = self.heartbeat_interval;
config.election_timeout_min = self.election_timeout_min;
config.election_timeout_max = self.election_timeout_max;
config.install_snapshot_timeout = self.install_snapshot_timeout;
}
}
#[derive(Clone, Debug)]
pub struct MetaMembership {
pub meta_group_id: RaftGroupId,
pub bootstrap_voters: Option<Vec<(NodeId, String)>>,
}
#[derive(Clone, Debug)]
pub struct NodeRuntimeConfig {
pub node_data: PathBuf,
pub security: TransportSecurity,
pub transport: TransportConfig,
pub listen_address: String,
pub rpc_address: Option<String>,
pub peers: Vec<(NodeId, String)>,
pub meta: Option<MetaMembership>,
pub timing: Option<GroupTiming>,
}
impl NodeRuntimeConfig {
pub fn new(node_data: PathBuf, listen_address: String) -> Self {
Self {
node_data,
security: TransportSecurity::PlaintextForTesting,
transport: TransportConfig::default(),
listen_address,
rpc_address: None,
peers: Vec::new(),
meta: None,
timing: None,
}
}
}
#[derive(Clone, Debug)]
pub struct MetaGroupStatus {
pub meta_group_id: RaftGroupId,
pub metadata_version: MetadataVersion,
pub metrics: GroupMetrics,
}
#[derive(Clone, Debug)]
pub struct TabletGroupStatus {
pub tablet_id: TabletId,
pub raft_group_id: RaftGroupId,
pub state: TabletState,
pub replicas: Vec<ReplicaDescriptor>,
pub applied: LogPosition,
pub metrics: GroupMetrics,
}
#[derive(Clone, Debug)]
pub struct RuntimeStatus {
pub identity: NodeIdentity,
pub rpc_address: String,
pub meta: Option<MetaGroupStatus>,
pub tablets: Vec<TabletGroupStatus>,
}
fn resolve_tablet_database_id(descriptor: &crate::tablet::TabletDescriptor) -> DatabaseId {
if descriptor.database_id != DatabaseId::ZERO {
return descriptor.database_id;
}
DatabaseId::from_bytes(*descriptor.raft_group_id.as_bytes())
}
fn tablet_group_name(tablet_id: TabletId) -> String {
format!("tablet-{}", tablet_id.to_hex())
}
pub const TABLET_LEDGER_FILENAME: &str = "tablet-ledger.json";
#[cfg(test)]
pub const TABLET_LEDGER_FORMAT_VERSION: u32 = 1;
#[cfg(test)]
pub const MIN_SUPPORTED_TABLET_LEDGER_FORMAT_VERSION: u32 = 1;
#[cfg(test)]
const MAX_TIMESTAMP: HlcTimestamp = HlcTimestamp {
physical_micros: u64::MAX,
logical: u32::MAX,
node_tiebreaker: u32::MAX,
};
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[cfg(test)]
struct TabletLedgerCheckpoint {
format_version: u32,
position: LogPosition,
rows: BTreeMap<Key, Vec<(HlcTimestamp, Vec<u8>)>>,
}
#[derive(Debug)]
#[cfg(test)]
pub struct TabletLedger {
rows: BTreeMap<Key, Vec<(HlcTimestamp, Vec<u8>)>>,
pins: BTreeMap<HlcTimestamp, usize>,
position: LogPosition,
state_dir: PathBuf,
}
#[cfg(test)]
impl TabletLedger {
pub fn open(group_dir: &Path) -> Result<Self, RuntimeError> {
let state_dir = group_dir.join("raft").join("state");
std::fs::create_dir_all(&state_dir).map_err(RuntimeError::Io)?;
let path = state_dir.join(TABLET_LEDGER_FILENAME);
let Some(bytes) = crate::node::read_meta_file(&path)? else {
return Ok(Self {
rows: BTreeMap::new(),
pins: BTreeMap::new(),
position: LogPosition::ZERO,
state_dir,
});
};
let checkpoint: TabletLedgerCheckpoint =
crate::node::decode_json(TABLET_LEDGER_FILENAME, &bytes)?;
if checkpoint.format_version < MIN_SUPPORTED_TABLET_LEDGER_FORMAT_VERSION
|| checkpoint.format_version > TABLET_LEDGER_FORMAT_VERSION
{
return Err(ClusterError::UnsupportedFormatVersion {
file: TABLET_LEDGER_FILENAME,
found: checkpoint.format_version,
min: MIN_SUPPORTED_TABLET_LEDGER_FORMAT_VERSION,
max: TABLET_LEDGER_FORMAT_VERSION,
}
.into());
}
Ok(Self {
rows: checkpoint.rows,
pins: BTreeMap::new(),
position: checkpoint.position,
state_dir,
})
}
fn apply(
&mut self,
command: &TabletDataCommand,
commit_ts: HlcTimestamp,
position: LogPosition,
) -> Result<(), RuntimeError> {
if position.index <= self.position.index {
return Ok(());
}
match command {
TabletDataCommand::Upsert { entries } => {
for (key, value) in entries {
self.insert_version(Key::from_bytes(key.clone()), commit_ts, value.clone());
}
}
TabletDataCommand::Delete { keys } => {
if !self.pins.is_empty() {
return Err(RuntimeError::InvalidRequest(
"legacy tablet ledger Delete with live snapshot pins".to_owned(),
));
}
for key in keys {
self.rows.remove(&Key::from_bytes(key.clone()));
}
}
TabletDataCommand::Replace { rows } => {
if !self.pins.is_empty() {
return Err(RuntimeError::InvalidRequest(
"tablet ledger Replace with live snapshot pins".to_owned(),
));
}
self.rows.clear();
for (key, value) in rows {
self.insert_version(Key::from_bytes(key.clone()), commit_ts, value.clone());
}
}
}
self.position = position;
self.persist()
}
fn insert_version(&mut self, key: Key, ts: HlcTimestamp, value: Vec<u8>) {
let chain = self.rows.entry(key).or_default();
chain.push((ts, value));
chain.sort_by_key(|(version, _)| *version);
match self.pins.keys().next() {
None => {
let newest = chain.pop().expect("just inserted");
chain.clear();
chain.push(newest);
}
Some(oldest_pin) => {
let baseline = chain.partition_point(|(version, _)| *version <= *oldest_pin);
if baseline > 1 {
chain.drain(..baseline - 1);
}
}
}
}
fn persist(&self) -> Result<(), RuntimeError> {
let checkpoint = TabletLedgerCheckpoint {
format_version: TABLET_LEDGER_FORMAT_VERSION,
position: self.position,
rows: self.rows.clone(),
};
let bytes = crate::node::encode_json(TABLET_LEDGER_FILENAME, &checkpoint)?;
crate::node::write_meta_atomic(&self.state_dir, TABLET_LEDGER_FILENAME, &bytes)
.map_err(ClusterError::Io)?;
Ok(())
}
pub fn applied_position(&self) -> LogPosition {
self.position
}
fn pin(&mut self, ts: HlcTimestamp) {
*self.pins.entry(ts).or_insert(0) += 1;
}
fn unpin(&mut self, ts: HlcTimestamp) {
if let Some(count) = self.pins.get_mut(&ts) {
*count -= 1;
if *count == 0 {
self.pins.remove(&ts);
}
}
}
pub fn pin_count(&self) -> usize {
self.pins.values().sum()
}
pub fn rows_at(&self, ts: HlcTimestamp) -> BTreeMap<Key, Vec<u8>> {
self.rows
.iter()
.filter_map(|(key, chain)| {
let visible = chain.iter().rfind(|(version, _)| *version <= ts)?;
Some((key.clone(), visible.1.clone()))
})
.collect()
}
pub fn current_rows(&self) -> BTreeMap<Key, Vec<u8>> {
self.rows_at(MAX_TIMESTAMP)
}
pub fn deltas_after(&self, ts: HlcTimestamp) -> Vec<(Key, Vec<u8>)> {
let mut deltas: Vec<(HlcTimestamp, Key, Vec<u8>)> = Vec::new();
for (key, chain) in &self.rows {
for (version, value) in chain {
if *version > ts {
deltas.push((*version, key.clone(), value.clone()));
}
}
}
deltas.sort_by(|left, right| left.0.cmp(&right.0).then_with(|| left.1.cmp(&right.1)));
deltas
.into_iter()
.map(|(_, key, value)| (key, value))
.collect()
}
pub fn size_bytes(&self) -> u64 {
self.rows
.iter()
.map(|(key, chain)| {
let newest = chain.last().expect("non-empty chain");
(key.as_bytes().len() + newest.1.len()) as u64
})
.sum()
}
fn snapshot_bytes(
&self,
) -> Result<Vec<u8>, mongreldb_consensus::state_machine::StateMachineError> {
serde_json::to_vec(&TabletLedgerCheckpoint {
format_version: TABLET_LEDGER_FORMAT_VERSION,
position: self.position,
rows: self.rows.clone(),
})
.map_err(|error| {
mongreldb_consensus::state_machine::StateMachineError::Sink(format!(
"tablet ledger snapshot: {error}"
))
})
}
fn install_bytes(
&mut self,
bytes: &[u8],
) -> Result<(), mongreldb_consensus::state_machine::StateMachineError> {
use mongreldb_consensus::state_machine::StateMachineError;
let checkpoint: TabletLedgerCheckpoint =
serde_json::from_slice(bytes).map_err(|error| {
StateMachineError::Corrupt(format!("tablet ledger snapshot: {error}"))
})?;
if checkpoint.format_version < MIN_SUPPORTED_TABLET_LEDGER_FORMAT_VERSION
|| checkpoint.format_version > TABLET_LEDGER_FORMAT_VERSION
{
return Err(StateMachineError::Corrupt(format!(
"tablet ledger snapshot format version {} is outside \
{MIN_SUPPORTED_TABLET_LEDGER_FORMAT_VERSION}..={TABLET_LEDGER_FORMAT_VERSION}",
checkpoint.format_version
)));
}
self.rows = checkpoint.rows;
self.position = checkpoint.position;
self.persist()
.map_err(|error| StateMachineError::Sink(error.to_string()))
}
}
#[cfg(test)]
fn encode_group_snapshot(engine: &[u8], ledger: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(12 + engine.len() + ledger.len());
out.extend_from_slice(&1_u32.to_le_bytes());
out.extend_from_slice(&(engine.len() as u64).to_le_bytes());
out.extend_from_slice(engine);
out.extend_from_slice(ledger);
out
}
#[cfg(test)]
fn decode_group_snapshot(
bytes: &[u8],
) -> Result<(Vec<u8>, Vec<u8>), mongreldb_consensus::state_machine::StateMachineError> {
use mongreldb_consensus::state_machine::StateMachineError;
if bytes.len() < 12 {
return Err(StateMachineError::Corrupt(
"tablet group snapshot: truncated frame".to_owned(),
));
}
let version = u32::from_le_bytes(bytes[..4].try_into().expect("4 bytes"));
if version != 1 {
return Err(StateMachineError::Corrupt(format!(
"tablet group snapshot format version {version} is not 1"
)));
}
let engine_len = u64::from_le_bytes(bytes[4..12].try_into().expect("8 bytes")) as usize;
if bytes.len() < 12 + engine_len {
return Err(StateMachineError::Corrupt(
"tablet group snapshot: truncated engine payload".to_owned(),
));
}
Ok((
bytes[12..12 + engine_len].to_vec(),
bytes[12 + engine_len..].to_vec(),
))
}
struct TabletGroup {
group: Arc<ConsensusGroup<TcpTransport>>,
descriptor: TabletDescriptor,
sink: Arc<Mutex<EngineApplySink>>,
_ownership: TabletOwnershipGuard<'static>,
}
#[derive(serde::Deserialize)]
struct TabletFileProbe {
tablet: TabletProbe,
}
#[derive(serde::Deserialize)]
struct TabletProbe {
raft_group_id: RaftGroupId,
}
fn scan_tablet_layouts(node_data: &Path) -> Result<Vec<(TabletId, RaftGroupId)>, RuntimeError> {
let root = node_data.join(TABLETS_DIR);
if !root.is_dir() {
return Ok(Vec::new());
}
let mut found = Vec::new();
for entry in std::fs::read_dir(&root)? {
let entry = entry?;
if !entry.file_type()?.is_dir() {
continue;
}
let name = entry.file_name();
let Some(name) = name.to_str() else {
return Err(RuntimeError::InvalidRequest(format!(
"tablet directory name {name:?} is not UTF-8"
)));
};
let tablet_id: TabletId = name.parse().map_err(|_| {
RuntimeError::InvalidRequest(format!("tablet directory `{name}` is not a tablet id"))
})?;
let probe_path = entry.path().join(TABLET_META_FILENAME);
let Some(bytes) = crate::node::read_meta_file(&probe_path)? else {
return Err(TabletError::MissingMetadata(probe_path).into());
};
let probe: TabletFileProbe = serde_json::from_slice(&bytes).map_err(|error| {
RuntimeError::InvalidRequest(format!(
"tablet metadata probe at {} failed: {error}",
probe_path.display()
))
})?;
found.push((tablet_id, probe.tablet.raft_group_id));
}
found.sort();
Ok(found)
}
fn active_snapshot_timestamp(
node_data: &Path,
tablet_id: TabletId,
) -> Result<Option<HlcTimestamp>, RuntimeError> {
let mut found = None;
for (source_id, raft_group_id) in scan_tablet_layouts(node_data)? {
let layout = TabletLayout::new(node_data.to_path_buf(), source_id, raft_group_id);
let timestamp = if let Some(progress) = split_progress(&layout)?.filter(|progress| {
progress.source.tablet_id == tablet_id && progress.phase < SplitPhase::Published
}) {
Some(progress.split_ts)
} else {
merge_progress(&layout)?
.filter(|progress| {
progress.phase < MergePhase::Published
&& progress
.sources
.iter()
.any(|source| source.tablet_id == tablet_id)
})
.map(|progress| progress.merge_ts)
};
if let Some(timestamp) = timestamp {
if found
.replace(timestamp)
.is_some_and(|prior| prior != timestamp)
{
return Err(RuntimeError::InvalidRequest(format!(
"tablet {tablet_id} appears in conflicting split/merge progress"
)));
}
}
}
Ok(found)
}
fn peer_endpoint(security: &TransportSecurity, node_id: NodeId, address: &str) -> PeerEndpoint {
match security {
TransportSecurity::Mtls(_) => PeerEndpoint::mtls(address, node_id),
TransportSecurity::PlaintextForTesting => PeerEndpoint::plaintext(address),
}
}
struct TabletOpenContext<'a> {
node_data: &'a Path,
security: &'a TransportSecurity,
timing: Option<GroupTiming>,
identity: &'a NodeIdentity,
peers: &'a BTreeMap<NodeId, String>,
}
pub struct NodeRuntime {
identity: NodeIdentity,
node_data: PathBuf,
rpc_address: String,
security: TransportSecurity,
peers: BTreeMap<NodeId, String>,
timing: Option<GroupTiming>,
transport: Arc<TcpTransport>,
server: Option<TransportServer>,
meta: Option<Arc<MetaGroup<TcpTransport>>>,
tablets: BTreeMap<TabletId, TabletGroup>,
}
#[derive(Clone)]
pub struct NodeInternalRpcClient {
transport: Arc<TcpTransport>,
}
impl NodeInternalRpcClient {
pub async fn call(
&self,
target: NodeId,
service_id: u32,
body: Vec<u8>,
) -> Result<Vec<u8>, RuntimeError> {
self.transport
.internal_rpc(raft_node_id(&target), service_id, body)
.await
.map_err(Into::into)
}
}
impl NodeRuntime {
pub async fn start(config: NodeRuntimeConfig) -> Result<Self, RuntimeError> {
let identity =
NodeIdentity::load(&config.node_data)?.ok_or(ClusterError::NotInitialized)?;
let transport = Arc::new(TcpTransport::new(
config.transport.clone(),
config.security.clone(),
));
let mut peers = BTreeMap::new();
for (node_id, address) in &config.peers {
transport.upsert_peer(
raft_node_id(node_id),
peer_endpoint(&config.security, *node_id, address),
);
peers.insert(*node_id, address.clone());
}
let server = TransportServer::bind_shared(
&config.listen_address,
transport.security_handle(),
transport.registry(),
config.transport.clone(),
)
.await?;
let rpc_address = config
.rpc_address
.clone()
.unwrap_or_else(|| server.local_addr().to_string());
let meta = match &config.meta {
Some(membership) => {
let group =
Self::start_meta_group(&config, &identity, membership, transport.clone())
.await?;
Some(Arc::new(group))
}
None => None,
};
let tablets =
Self::open_tablet_groups(&config, &identity, &peers, transport.clone()).await?;
Ok(Self {
identity,
node_data: config.node_data.clone(),
rpc_address,
security: config.security.clone(),
peers,
timing: config.timing,
transport,
server: Some(server),
meta,
tablets,
})
}
async fn start_meta_group(
config: &NodeRuntimeConfig,
identity: &NodeIdentity,
membership: &MetaMembership,
transport: Arc<TcpTransport>,
) -> Result<MetaGroup<TcpTransport>, RuntimeError> {
let meta_config = MetaGroupConfig::new(
config.node_data.clone(),
membership.meta_group_id,
identity.node_id,
);
let mut group_config = meta_config.group_config();
if let Some(timing) = &config.timing {
timing.apply(&mut group_config);
}
let group = MetaGroup::create(meta_config, group_config, transport).await?;
if let Some(voters) = &membership.bootstrap_voters {
if !voters
.iter()
.any(|(node_id, _)| *node_id == identity.node_id)
{
return Err(RuntimeError::InvalidRequest(
"meta bootstrap voter set does not include this node".to_owned(),
));
}
if !group.is_initialized().await? {
group.bootstrap(voters).await?;
}
}
Ok(group)
}
async fn open_tablet_groups(
config: &NodeRuntimeConfig,
identity: &NodeIdentity,
peers: &BTreeMap<NodeId, String>,
transport: Arc<TcpTransport>,
) -> Result<BTreeMap<TabletId, TabletGroup>, RuntimeError> {
let mut groups = BTreeMap::new();
let context = TabletOpenContext {
node_data: &config.node_data,
security: &config.security,
timing: config.timing,
identity,
peers,
};
for (tablet_id, raft_group_id) in scan_tablet_layouts(&config.node_data)? {
let layout = TabletLayout::new(config.node_data.clone(), tablet_id, raft_group_id);
let descriptor = layout.validate()?;
let group =
Self::open_tablet_group(&context, transport.clone(), &layout, &descriptor).await?;
groups.insert(tablet_id, group);
}
Ok(groups)
}
async fn open_tablet_group(
context: &TabletOpenContext<'_>,
transport: Arc<TcpTransport>,
layout: &TabletLayout,
descriptor: &TabletDescriptor,
) -> Result<TabletGroup, RuntimeError> {
let replica = *descriptor
.replica_on(context.identity.node_id)
.ok_or_else(|| {
RuntimeError::InvalidRequest(format!(
"tablet {} descriptor does not list this node ({}) as a replica",
descriptor.tablet_id, context.identity.node_id
))
})?;
if transport.registry().get(replica.raft_node_id).is_some() {
return Err(RuntimeError::InvalidRequest(format!(
"raft id {} is already attached to this node's transport registry",
replica.raft_node_id
)));
}
for other in &descriptor.replicas {
let Some(address) = context.peers.get(&other.node_id) else {
return Err(RuntimeError::InvalidRequest(format!(
"no membership-directory address for replica node {} of tablet {}",
other.node_id, descriptor.tablet_id
)));
};
transport.upsert_peer(
other.raft_node_id,
peer_endpoint(context.security, other.node_id, address),
);
}
let ownership = TabletOwnershipRegistry::global().try_reserve(layout)?;
let database_id = resolve_tablet_database_id(descriptor);
let engine_config = EngineGroupConfig::new(
context.node_data.to_path_buf(),
descriptor.raft_group_id,
context.identity.cluster_id,
context.identity.node_id,
database_id,
);
let sink = open_engine_sink(&engine_config)?;
{
let mut engine = sink
.lock()
.map_err(|_| RuntimeError::InvalidRequest("engine sink lock poisoned".into()))?;
engine.initialize_tablet_keyspace()?;
let legacy_path = engine_config
.group_dir()
.join("raft")
.join("state")
.join(TABLET_LEDGER_FILENAME);
if legacy_path.is_file() {
let bytes = std::fs::read(&legacy_path)?;
engine.migrate_legacy_tablet_ledger(
&bytes,
active_snapshot_timestamp(context.node_data, descriptor.tablet_id)?,
)?;
std::fs::remove_file(&legacy_path)?;
if let Some(parent) = legacy_path.parent() {
std::fs::File::open(parent)?.sync_all()?;
}
}
engine.finish_legacy_tablet_migration()?;
}
let mut group_config = GroupConfig::new(
tablet_group_name(descriptor.tablet_id),
replica.raft_node_id,
engine_config.group_dir(),
);
group_config.storage = engine_config.storage.clone();
group_config.idempotency_retention = engine_config.idempotency_retention;
if let Some(timing) = &context.timing {
timing.apply(&mut group_config);
}
let dyn_sink: Arc<Mutex<dyn ApplySink>> = sink.clone();
let group = Arc::new(ConsensusGroup::create(group_config, transport, dyn_sink).await?);
let commit_log: Arc<dyn CommitLog> = Arc::new(RaftCommitLog::new(group.clone()));
sink.lock()
.map_err(|_| RuntimeError::InvalidRequest("engine sink lock poisoned".into()))?
.bind_commit_log(commit_log)?;
Ok(TabletGroup {
group,
descriptor: descriptor.clone(),
sink,
_ownership: ownership,
})
}
pub fn identity(&self) -> &NodeIdentity {
&self.identity
}
pub fn reload_trust(
&mut self,
trust: crate::bootstrap::TrustConfig,
) -> Result<(), RuntimeError> {
let tls = crate::network::TlsConfig::from_trust(&trust)
.map_err(|e| RuntimeError::InvalidRequest(format!("reload trust: {e}")))?;
let security = crate::network::TransportSecurity::Mtls(tls);
self.transport.reload_security(security.clone());
self.security = security;
Ok(())
}
pub fn rpc_address(&self) -> &str {
&self.rpc_address
}
pub fn attach_internal_rpc_handler(
&self,
service_id: u32,
handler: Arc<dyn InternalRpcHandler>,
) {
self.transport.registry().attach_internal(
raft_node_id(&self.identity.node_id),
service_id,
handler,
);
}
pub async fn internal_rpc(
&self,
target: NodeId,
service_id: u32,
body: Vec<u8>,
) -> Result<Vec<u8>, RuntimeError> {
self.transport
.internal_rpc(raft_node_id(&target), service_id, body)
.await
.map_err(Into::into)
}
pub fn internal_rpc_client(&self) -> NodeInternalRpcClient {
NodeInternalRpcClient {
transport: Arc::clone(&self.transport),
}
}
pub fn meta_group(&self) -> Option<&MetaGroup<TcpTransport>> {
self.meta.as_deref()
}
pub fn tablet_group(&self, tablet_id: TabletId) -> Option<&ConsensusGroup<TcpTransport>> {
self.tablets.get(&tablet_id).map(|tablet| &*tablet.group)
}
fn tablet_sink(&self, tablet_id: TabletId) -> Option<Arc<Mutex<EngineApplySink>>> {
self.tablets
.get(&tablet_id)
.map(|tablet| tablet.sink.clone())
}
pub fn tablet_descriptor(&self, tablet_id: TabletId) -> Option<&TabletDescriptor> {
self.tablets
.get(&tablet_id)
.map(|tablet| &tablet.descriptor)
}
pub fn tablet_ids(&self) -> Vec<TabletId> {
self.tablets.keys().copied().collect()
}
fn hosted_tablet(&self, tablet_id: TabletId) -> Result<&TabletGroup, RuntimeError> {
self.tablets.get(&tablet_id).ok_or_else(|| {
RuntimeError::InvalidRequest(format!("this node hosts no tablet {tablet_id}"))
})
}
pub async fn create_tablet(
&mut self,
descriptor: &TabletDescriptor,
partitioning: &TablePartitioningRecord,
bootstrap_voters: Option<&[(NodeId, String)]>,
publish_to_meta: bool,
control: &ExecutionControl,
) -> Result<(), RuntimeError> {
descriptor.validate()?;
partitioning.validate()?;
if partitioning.table_id != descriptor.table_id {
return Err(RuntimeError::InvalidRequest(format!(
"partitioning record names table {} but the descriptor partitions table {}",
partitioning.table_id, descriptor.table_id
)));
}
self.create_hosted_replica(descriptor, bootstrap_voters)
.await?;
if publish_to_meta {
self.publish_tablet_descriptor(descriptor, control).await?;
}
Ok(())
}
async fn create_hosted_replica(
&mut self,
descriptor: &TabletDescriptor,
bootstrap_voters: Option<&[(NodeId, String)]>,
) -> Result<(), RuntimeError> {
descriptor.validate()?;
if let Some(existing) = self.tablets.get(&descriptor.tablet_id) {
if existing.descriptor == *descriptor {
return Ok(());
}
return Err(RuntimeError::InvalidRequest(format!(
"tablet {} is already open with a different descriptor",
descriptor.tablet_id
)));
}
if descriptor.replica_on(self.identity.node_id).is_none() {
return Err(RuntimeError::InvalidRequest(format!(
"tablet {} descriptor does not list this node ({}) as a replica",
descriptor.tablet_id, self.identity.node_id
)));
}
let layout = TabletLayout::new(
self.node_data.clone(),
descriptor.tablet_id,
descriptor.raft_group_id,
);
layout.create(descriptor)?;
let context = TabletOpenContext {
node_data: &self.node_data,
security: &self.security,
timing: self.timing,
identity: &self.identity,
peers: &self.peers,
};
let group =
Self::open_tablet_group(&context, self.transport.clone(), &layout, descriptor).await?;
if let Some(voters) = bootstrap_voters {
let mut members = BTreeMap::new();
for (node_id, address) in voters {
let Some(replica) = descriptor.replica_on(*node_id) else {
return Err(RuntimeError::InvalidRequest(format!(
"bootstrap voter {node_id} is not a replica of tablet {}",
descriptor.tablet_id
)));
};
members.insert(replica.raft_node_id, BasicNode::new(address.clone()));
}
if !members.contains_key(&group.group.node_id()) {
return Err(RuntimeError::InvalidRequest(
"tablet bootstrap voter set does not include this node".to_owned(),
));
}
if !group.group.is_initialized().await? {
group.group.bootstrap(members).await?;
}
}
self.tablets.insert(descriptor.tablet_id, group);
Ok(())
}
pub async fn add_tablet_replica(
&self,
tablet_id: TabletId,
replica: ReplicaDescriptor,
address: &str,
) -> Result<(), RuntimeError> {
let tablet = self.hosted_tablet(tablet_id)?;
if tablet.descriptor.replicas.iter().any(|existing| {
existing.node_id == replica.node_id || existing.raft_node_id == replica.raft_node_id
}) {
return Err(RuntimeError::InvalidRequest(format!(
"replica node {} / raft id {} already belongs to tablet {}",
replica.node_id, replica.raft_node_id, tablet_id
)));
}
self.transport.upsert_peer(
replica.raft_node_id,
peer_endpoint(&self.security, replica.node_id, address),
);
for _ in 0..10 {
tablet
.group
.wait_uniform_membership(Duration::from_secs(30))
.await?;
let (voters, learners) = tablet.group.members();
if voters.contains(&replica.raft_node_id) {
return Ok(());
}
if !learners.contains(&replica.raft_node_id) {
match tablet
.group
.add_learner(replica.raft_node_id, BasicNode::new(address.to_owned()))
.await
{
Ok(()) => {}
Err(ConsensusError::MembershipInProgress) => continue,
Err(error) => return Err(error.into()),
}
}
match tablet.group.promote(replica.raft_node_id).await {
Ok(()) => return Ok(()),
Err(ConsensusError::MembershipInProgress) => continue,
Err(error) => return Err(error.into()),
}
}
Err(RuntimeError::InvalidRequest(format!(
"membership change for tablet {tablet_id} did not settle"
)))
}
pub async fn publish_tablet_descriptor(
&self,
descriptor: &TabletDescriptor,
control: &ExecutionControl,
) -> Result<MetadataVersion, RuntimeError> {
let meta = self.meta.as_ref().ok_or_else(|| {
RuntimeError::InvalidRequest("this node hosts no meta group".to_owned())
})?;
let receipt = meta
.propose(
crate::meta::new_command_id()?,
MetaCommand::SetTabletDescriptor {
descriptor: descriptor.clone(),
},
control,
)
.await?;
Ok(receipt.metadata_version)
}
fn meta_group_or_err(&self) -> Result<Arc<MetaGroup<TcpTransport>>, RuntimeError> {
self.meta
.clone()
.ok_or_else(|| RuntimeError::InvalidRequest("this node hosts no meta group".to_owned()))
}
async fn plan_split(
&self,
tablet_id: TabletId,
split_key: Option<Key>,
control: &ExecutionControl,
) -> Result<SplitPlan, RuntimeError> {
let meta = self.meta_group_or_err()?;
let source = meta.state().tablet(tablet_id).cloned().ok_or_else(|| {
RuntimeError::InvalidRequest(format!("tablet {tablet_id} is not in the meta state"))
})?;
if source.state != TabletState::Active {
return Err(RuntimeError::InvalidRequest(format!(
"tablet {tablet_id} is in state {}, expected Active",
source.state
)));
}
if source.replica_on(self.identity.node_id).is_none() {
return Err(RuntimeError::InvalidRequest(format!(
"this node hosts no replica of tablet {tablet_id}"
)));
}
let replica_count = source.replicas.len();
let raft_ids = meta
.allocate_raft_node_ids(
2 * u32::try_from(replica_count).unwrap_or(u32::MAX),
control,
)
.await?;
let allocation = |ids: &[RaftNodeId]| ChildAllocation {
tablet_id: TabletId::new_random(),
raft_group_id: RaftGroupId::new_random(),
replicas: source
.replicas
.iter()
.zip(ids)
.map(|(replica, raft_node_id)| ReplicaDescriptor {
node_id: replica.node_id,
role: replica.role,
raft_node_id: *raft_node_id,
})
.collect(),
};
let selection = match split_key {
Some(key) => SplitKeySelection::Explicit(key),
None => SplitKeySelection::Midpoint,
};
let plan = TabletSplitPlanner::new(self.node_data.clone()).plan(
&source,
selection,
now_timestamp(),
[
allocation(&raft_ids[..replica_count]),
allocation(&raft_ids[replica_count..]),
],
)?;
Ok(plan)
}
async fn child_sinks(
&mut self,
children: &[TabletDescriptor; 2],
control: &ExecutionControl,
) -> Result<[RuntimeChildSink; 2], RuntimeError> {
for child in children {
if child.replica_on(self.identity.node_id).is_some()
&& !self.tablets.contains_key(&child.tablet_id)
{
let voters: Vec<(NodeId, String)> = child
.replicas
.iter()
.filter_map(|replica| {
self.peers
.get(&replica.node_id)
.map(|address| (replica.node_id, address.clone()))
})
.collect();
self.create_hosted_replica(child, Some(voters.as_slice()))
.await?;
}
}
let mut sinks = Vec::with_capacity(2);
for child in children {
let group = self
.tablets
.get(&child.tablet_id)
.ok_or_else(|| {
RuntimeError::InvalidRequest(format!(
"child tablet {} is not hosted on this node",
child.tablet_id
))
})?
.group
.clone();
sinks.push(RuntimeChildSink {
group,
handle: tokio::runtime::Handle::current(),
control: control.clone(),
staged: None,
});
}
let [lower, upper]: [RuntimeChildSink; 2] = sinks.try_into().map_err(|_| {
RuntimeError::InvalidRequest("expected exactly two child sinks".to_owned())
})?;
Ok([lower, upper])
}
async fn drop_hosted_tablet(&mut self, tablet_id: TabletId) -> Result<(), RuntimeError> {
if let Some(tablet) = self.tablets.remove(&tablet_id) {
tablet.group.shutdown().await?;
}
Ok(())
}
fn refresh_local_descriptor(&mut self, layout: &TabletLayout) -> Result<(), RuntimeError> {
let descriptor = layout.load_metadata()?;
if let Some(tablet) = self.tablets.get_mut(&layout.tablet_id()) {
if descriptor.generation > tablet.descriptor.generation {
tablet.descriptor = descriptor;
}
}
Ok(())
}
pub async fn split_step(
&mut self,
tablet_id: TabletId,
split_key: Option<Key>,
control: &ExecutionControl,
) -> Result<(SplitPhase, Option<SplitPublishCommand>), RuntimeError> {
let meta = self.meta_group_or_err()?;
let raft_group_id = match self.tablets.get(&tablet_id) {
Some(tablet) => tablet.descriptor.raft_group_id,
None => meta
.state()
.tablet(tablet_id)
.map(|descriptor| descriptor.raft_group_id)
.ok_or_else(|| {
RuntimeError::InvalidRequest(format!("this node hosts no tablet {tablet_id}"))
})?,
};
let source_layout = TabletLayout::new(self.node_data.clone(), tablet_id, raft_group_id);
let progress = split_progress(&source_layout)?;
let hosted = self.tablets.contains_key(&tablet_id);
if !hosted
&& progress
.as_ref()
.is_none_or(|record| record.phase < SplitPhase::Published)
{
return Err(RuntimeError::InvalidRequest(format!(
"this node hosts no tablet {tablet_id}"
)));
}
let keyspace = RuntimeKeyspace {
sink: if progress
.as_ref()
.is_some_and(|record| record.phase >= SplitPhase::Published)
{
None
} else {
Some(self.tablet_sink(tablet_id).ok_or_else(|| {
RuntimeError::InvalidRequest(format!("this node hosts no tablet {tablet_id}"))
})?)
},
pinned: None,
};
let plane = RuntimeMetaPlane {
meta: meta.clone(),
handle: tokio::runtime::Handle::current(),
control: control.clone(),
};
let mut executor = match progress {
Some(progress) => {
let children = progress.plan().child_descriptors();
let sinks = self.child_sinks(&children, control).await?;
SplitExecutor::resume(source_layout.clone(), plane, keyspace, sinks)?.ok_or_else(
|| {
RuntimeError::InvalidRequest(format!(
"split progress of tablet {tablet_id} vanished mid-resume"
))
},
)?
}
None => {
let plan = self.plan_split(tablet_id, split_key, control).await?;
let sinks = self.child_sinks(&plan.child_descriptors(), control).await?;
SplitExecutor::begin(plan, source_layout.clone(), plane, keyspace, sinks)?
}
};
if executor.phase() == SplitPhase::Published {
if let Some(sink) = self.tablet_sink(tablet_id) {
sink.lock()
.map_err(|_| RuntimeError::InvalidRequest("engine sink lock poisoned".into()))?
.release_tablet_snapshot()?;
}
self.drop_hosted_tablet(tablet_id).await?;
}
let phase = tokio::task::spawn_blocking(move || executor.step())
.await
.map_err(|error| {
RuntimeError::InvalidRequest(format!("split driver task failed: {error}"))
})??;
if phase == SplitPhase::Published {
if let Some(sink) = self.tablet_sink(tablet_id) {
sink.lock()
.map_err(|_| RuntimeError::InvalidRequest("engine sink lock poisoned".into()))?
.release_tablet_snapshot()?;
}
}
if matches!(phase, SplitPhase::MarkedSplitting | SplitPhase::Published) {
self.refresh_local_descriptor(&source_layout)?;
}
let published = if phase >= SplitPhase::Published {
match split_progress(&source_layout)? {
Some(record) => Some(SplitPublishCommand::from_plan(&record.plan())?),
None => None,
}
} else {
None
};
Ok((phase, published))
}
pub async fn split_tablet(
&mut self,
tablet_id: TabletId,
split_key: Option<Key>,
control: &ExecutionControl,
) -> Result<SplitPublishCommand, RuntimeError> {
let mut published = None;
loop {
let (phase, command) = self
.split_step(tablet_id, split_key.clone(), control)
.await?;
if command.is_some() {
published = command;
}
if phase == SplitPhase::SourceRetired {
break;
}
}
published.ok_or_else(|| {
RuntimeError::InvalidRequest(format!(
"split of tablet {tablet_id} completed without a publication"
))
})
}
pub async fn abort_split(
&mut self,
tablet_id: TabletId,
control: &ExecutionControl,
) -> Result<SplitAbortReport, RuntimeError> {
let meta = self.meta_group_or_err()?;
let source_layout = {
let source = self.hosted_tablet(tablet_id)?;
TabletLayout::new(
self.node_data.clone(),
tablet_id,
source.descriptor.raft_group_id,
)
};
if let Some(progress) = split_progress(&source_layout)? {
if progress.phase >= SplitPhase::Published {
return Err(SplitError::CannotAbort {
tablet: tablet_id,
phase: progress.phase,
}
.into());
}
for child in &progress.children {
self.drop_hosted_tablet(child.tablet_id).await?;
}
}
let mut plane = RuntimeMetaPlane {
meta,
handle: tokio::runtime::Handle::current(),
control: control.clone(),
};
let layout = source_layout.clone();
let report = tokio::task::spawn_blocking(move || abort_split(&layout, &mut plane))
.await
.map_err(|error| {
RuntimeError::InvalidRequest(format!("split abort task failed: {error}"))
})??;
if let Some(sink) = self.tablet_sink(tablet_id) {
sink.lock()
.map_err(|_| RuntimeError::InvalidRequest("engine sink lock poisoned".into()))?
.release_tablet_snapshot()?;
}
self.refresh_local_descriptor(&source_layout)?;
Ok(report)
}
pub async fn merge_step(
&mut self,
first: TabletId,
second: TabletId,
control: &ExecutionControl,
) -> Result<(MergePhase, Option<MergePublishCommand>), RuntimeError> {
let meta = self.meta_group_or_err()?;
let layout_of = |tablet_id: TabletId| -> Result<TabletLayout, RuntimeError> {
let raft_group_id = match self.tablets.get(&tablet_id) {
Some(tablet) => tablet.descriptor.raft_group_id,
None => meta
.state()
.tablet(tablet_id)
.map(|descriptor| descriptor.raft_group_id)
.ok_or_else(|| {
RuntimeError::InvalidRequest(format!(
"this node hosts no tablet {tablet_id}"
))
})?,
};
Ok(TabletLayout::new(
self.node_data.clone(),
tablet_id,
raft_group_id,
))
};
let first_layout = layout_of(first)?;
let second_layout = layout_of(second)?;
let phase_in_flight = match (
merge_progress(&first_layout)?,
merge_progress(&second_layout)?,
) {
(Some(progress), None) => Some(progress.phase),
(None, None) => None,
_ => {
return Err(RuntimeError::InvalidRequest(
"merge progress records on both sources: corrupt state".to_owned(),
));
}
};
let retired_only = matches!(
phase_in_flight,
Some(MergePhase::Published | MergePhase::SourcesRetired)
);
if !retired_only
&& (!self.tablets.contains_key(&first) || !self.tablets.contains_key(&second))
{
return Err(RuntimeError::InvalidRequest(format!(
"this node must host both merge sources ({first}, {second})"
)));
}
let keyspace_of = |layout: &TabletLayout| -> Result<RuntimeKeyspace, RuntimeError> {
Ok(RuntimeKeyspace {
sink: if retired_only {
None
} else {
Some(self.tablet_sink(layout.tablet_id()).ok_or_else(|| {
RuntimeError::InvalidRequest(format!(
"this node hosts no tablet {}",
layout.tablet_id()
))
})?)
},
pinned: None,
})
};
let (first_keyspace, second_keyspace) =
(keyspace_of(&first_layout)?, keyspace_of(&second_layout)?);
let ordered_layouts = |lower_first: bool| {
if lower_first {
[first_layout.clone(), second_layout.clone()]
} else {
[second_layout.clone(), first_layout.clone()]
}
};
let plane = RuntimeMetaPlane {
meta: meta.clone(),
handle: tokio::runtime::Handle::current(),
control: control.clone(),
};
let mut executor = match phase_in_flight {
Some(_) => {
let progress = merge_progress(&first_layout)?
.or(merge_progress(&second_layout)?)
.expect("phase_in_flight is Some");
let lower_first = progress.sources[0].tablet_id == first;
let sink = self.replacement_sink(&progress.plan(), control).await?;
let keyspaces = if lower_first {
[first_keyspace, second_keyspace]
} else {
[second_keyspace, first_keyspace]
};
MergeExecutor::resume(ordered_layouts(lower_first), plane, keyspaces, sink)?
.ok_or_else(|| {
RuntimeError::InvalidRequest(
"merge progress vanished mid-resume".to_owned(),
)
})?
}
None => {
let plan = self.plan_merge(first, second, control).await?;
let lower_first = plan.sources[0].tablet_id == first;
let sink = self.replacement_sink(&plan, control).await?;
let keyspaces = if lower_first {
[first_keyspace, second_keyspace]
} else {
[second_keyspace, first_keyspace]
};
MergeExecutor::begin(plan, ordered_layouts(lower_first), plane, keyspaces, sink)?
}
};
if executor.phase() == MergePhase::Published {
for tablet_id in [first, second] {
if let Some(sink) = self.tablet_sink(tablet_id) {
sink.lock()
.map_err(|_| {
RuntimeError::InvalidRequest("engine sink lock poisoned".into())
})?
.release_tablet_snapshot()?;
}
}
self.drop_hosted_tablet(first).await?;
self.drop_hosted_tablet(second).await?;
}
let phase = tokio::task::spawn_blocking(move || executor.step())
.await
.map_err(|error| {
RuntimeError::InvalidRequest(format!("merge driver task failed: {error}"))
})??;
if phase == MergePhase::Published {
for tablet_id in [first, second] {
if let Some(sink) = self.tablet_sink(tablet_id) {
sink.lock()
.map_err(|_| {
RuntimeError::InvalidRequest("engine sink lock poisoned".into())
})?
.release_tablet_snapshot()?;
}
}
}
if matches!(phase, MergePhase::MarkedMerging | MergePhase::Published) {
self.refresh_local_descriptor(&first_layout)?;
self.refresh_local_descriptor(&second_layout)?;
}
let published = if phase >= MergePhase::Published {
let progress = merge_progress(&first_layout)?.or(merge_progress(&second_layout)?);
match progress {
Some(record) => Some(MergePublishCommand::from_plan(&record.plan())?),
None => None,
}
} else {
None
};
Ok((phase, published))
}
pub async fn merge_tablets(
&mut self,
first: TabletId,
second: TabletId,
control: &ExecutionControl,
) -> Result<MergePublishCommand, RuntimeError> {
let mut published = None;
loop {
let (phase, command) = self.merge_step(first, second, control).await?;
if command.is_some() {
published = command;
}
if phase == MergePhase::SourcesRetired {
break;
}
}
published.ok_or_else(|| {
RuntimeError::InvalidRequest("merge completed without a publication".to_owned())
})
}
async fn plan_merge(
&self,
first: TabletId,
second: TabletId,
control: &ExecutionControl,
) -> Result<MergePlan, RuntimeError> {
let meta = self.meta_group_or_err()?;
let state = meta.state();
let descriptor_of = |tablet_id: TabletId| -> Result<TabletDescriptor, RuntimeError> {
state.tablet(tablet_id).cloned().ok_or_else(|| {
RuntimeError::InvalidRequest(format!("tablet {tablet_id} is not in the meta state"))
})
};
let (first_desc, second_desc) = (descriptor_of(first)?, descriptor_of(second)?);
let schema = state.table(first_desc.table_id).ok_or_else(|| {
RuntimeError::InvalidRequest(format!(
"table {} of tablet {first} is not in the meta state",
first_desc.table_id
))
})?;
let active_schema_job = state
.schema_jobs
.values()
.find(|job| job.table_id == first_desc.table_id && !job.state.is_terminal())
.map(|job| job.job_id);
let size_of = |tablet_id: TabletId| -> Result<u64, RuntimeError> {
self.tablet_sink(tablet_id)
.ok_or_else(|| {
RuntimeError::InvalidRequest(format!("this node hosts no tablet {tablet_id}"))
})?
.lock()
.map_err(|_| RuntimeError::InvalidRequest("engine sink lock poisoned".into()))?
.tablet_size_bytes()
.map_err(RuntimeError::from)
};
let replica_count = first_desc.replicas.len();
let raft_ids = meta
.allocate_raft_node_ids(u32::try_from(replica_count).unwrap_or(u32::MAX), control)
.await?;
let allocation = ChildAllocation {
tablet_id: TabletId::new_random(),
raft_group_id: RaftGroupId::new_random(),
replicas: first_desc
.replicas
.iter()
.zip(&raft_ids)
.map(|(replica, raft_node_id)| ReplicaDescriptor {
node_id: replica.node_id,
role: replica.role,
raft_node_id: *raft_node_id,
})
.collect(),
};
let plan = MergePlanner::new(self.node_data.clone()).plan(
MergeInputs {
first: first_desc,
second: second_desc,
first_schema: schema.schema_version,
second_schema: schema.schema_version,
active_schema_job,
first_size_bytes: size_of(first)?,
second_size_bytes: size_of(second)?,
max_merged_size_bytes: DEFAULT_MAX_MERGED_SIZE_BYTES,
},
now_timestamp(),
allocation,
)?;
Ok(plan)
}
async fn replacement_sink(
&mut self,
plan: &MergePlan,
control: &ExecutionControl,
) -> Result<RuntimeChildSink, RuntimeError> {
let descriptor = plan.replacement_descriptor();
if descriptor.replica_on(self.identity.node_id).is_some()
&& !self.tablets.contains_key(&descriptor.tablet_id)
{
let voters: Vec<(NodeId, String)> = descriptor
.replicas
.iter()
.filter_map(|replica| {
self.peers
.get(&replica.node_id)
.map(|address| (replica.node_id, address.clone()))
})
.collect();
self.create_hosted_replica(&descriptor, Some(voters.as_slice()))
.await?;
}
let group = self
.tablets
.get(&descriptor.tablet_id)
.ok_or_else(|| {
RuntimeError::InvalidRequest(format!(
"replacement tablet {} is not hosted on this node",
descriptor.tablet_id
))
})?
.group
.clone();
Ok(RuntimeChildSink {
group,
handle: tokio::runtime::Handle::current(),
control: control.clone(),
staged: None,
})
}
pub async fn sync_hosted_tablets(
&mut self,
_control: &ExecutionControl,
) -> Result<HostedSyncReport, RuntimeError> {
let meta = self.meta_group_or_err()?;
let state = meta.state();
let node_id = self.identity.node_id;
let mut report = HostedSyncReport::default();
for record in state.tablets.values() {
let descriptor = &record.descriptor;
if self.tablets.contains_key(&descriptor.tablet_id)
|| descriptor.replica_on(node_id).is_none()
{
continue;
}
self.create_hosted_replica(descriptor, None).await?;
report.created.push(descriptor.tablet_id);
}
for record in state.tablets.values() {
let published = &record.descriptor;
let Some(tablet) = self.tablets.get(&published.tablet_id) else {
continue;
};
if published.generation <= tablet.descriptor.generation {
continue;
}
if published.raft_group_id != tablet.descriptor.raft_group_id {
return Err(RuntimeError::InvalidRequest(format!(
"meta descriptor for tablet {} names a different raft group than the \
hosted replica",
published.tablet_id
)));
}
if published.replica_on(node_id).is_none() {
continue; }
let layout = TabletLayout::new(
self.node_data.clone(),
published.tablet_id,
published.raft_group_id,
);
layout.store_metadata(published)?;
self.tablets
.get_mut(&published.tablet_id)
.expect("checked above")
.descriptor = published.clone();
report.refreshed.push(published.tablet_id);
}
let outgoing: Vec<TabletId> = self
.tablets
.keys()
.filter(|tablet_id| {
state
.tablets
.get(tablet_id)
.is_none_or(|record| record.descriptor.replica_on(node_id).is_none())
})
.copied()
.collect();
for tablet_id in outgoing {
let Some(tablet) = self.tablets.remove(&tablet_id) else {
continue;
};
let layout = TabletLayout::new(
self.node_data.clone(),
tablet_id,
tablet.descriptor.raft_group_id,
);
tablet.group.shutdown().await?;
drop(tablet);
layout.teardown()?;
report.torn_down.push(tablet_id);
}
Ok(report)
}
pub async fn write_tablet_rows(
&self,
tablet_id: TabletId,
entries: &[(Key, Vec<u8>)],
control: &ExecutionControl,
) -> Result<GroupCommitReceipt, RuntimeError> {
let tablet = self.hosted_tablet(tablet_id)?;
let envelope = CommandEnvelope::new(
COMMAND_TYPE_TABLET_DATA,
new_data_command_id()?,
TabletDataCommandRecord::new(TabletDataCommand::Upsert {
entries: entries
.iter()
.map(|(key, value)| (key.as_bytes().to_vec(), value.clone()))
.collect(),
})
.encode(),
);
let receipt = tablet
.group
.propose(CommandKind::Catalog, envelope, control)
.await?;
Ok(receipt)
}
pub fn tablet_rows(&self, tablet_id: TabletId) -> Result<BTreeMap<Key, Vec<u8>>, RuntimeError> {
let sink = self.tablet_sink(tablet_id).ok_or_else(|| {
RuntimeError::InvalidRequest(format!("this node hosts no tablet {tablet_id}"))
})?;
let rows = sink
.lock()
.map_err(|_| RuntimeError::InvalidRequest("engine sink lock poisoned".into()))?
.tablet_rows()?;
Ok(rows
.into_iter()
.map(|(key, value)| (Key::from_bytes(key), value))
.collect())
}
pub fn status(&self) -> RuntimeStatus {
RuntimeStatus {
identity: self.identity.clone(),
rpc_address: self.rpc_address.clone(),
meta: self.meta.as_ref().map(|meta| MetaGroupStatus {
meta_group_id: meta.meta_group_id(),
metadata_version: meta.metadata_version(),
metrics: meta.group().metrics(),
}),
tablets: self
.tablets
.values()
.map(|tablet| TabletGroupStatus {
tablet_id: tablet.descriptor.tablet_id,
raft_group_id: tablet.descriptor.raft_group_id,
state: tablet.descriptor.state,
replicas: tablet.descriptor.replicas.clone(),
applied: tablet.group.applied_position(),
metrics: tablet.group.metrics(),
})
.collect(),
}
}
pub async fn shutdown(mut self) -> Result<(), RuntimeError> {
if let Some(server) = self.server.take() {
server.shutdown().await;
}
let mut first_error: Option<RuntimeError> = None;
for (_, tablet) in std::mem::take(&mut self.tablets) {
if let Err(error) = tablet.group.shutdown().await {
if first_error.is_none() {
first_error = Some(error.into());
}
}
let close_result = tablet
.sink
.lock()
.map_err(|_| RuntimeError::InvalidRequest("engine sink lock poisoned".to_owned()))
.and_then(|mut sink| sink.close().map_err(RuntimeError::from));
if let Err(error) = close_result {
if first_error.is_none() {
first_error = Some(error);
}
}
}
if let Some(meta) = self.meta.take() {
if let Err(error) = meta.shutdown().await {
if first_error.is_none() {
first_error = Some(error.into());
}
}
}
match first_error {
Some(error) => Err(error),
None => Ok(()),
}
}
pub async fn crash(mut self) {
drop(self.server.take());
for (_, tablet) in std::mem::take(&mut self.tablets) {
let TabletGroup { group, sink, .. } = tablet;
match Arc::try_unwrap(group) {
Ok(group) => group.crash().await,
Err(group) => {
let _ = group.shutdown().await;
}
}
if let Ok(mut sink) = sink.lock() {
let _ = sink.crash();
};
}
if let Some(meta) = self.meta.take() {
match Arc::try_unwrap(meta) {
Ok(meta) => meta.crash().await,
Err(meta) => {
let _ = meta.shutdown().await;
}
}
}
}
}
pub const DEFAULT_MAX_MERGED_SIZE_BYTES: u64 = 64 * 1024 * 1024;
fn now_timestamp() -> HlcTimestamp {
let micros = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_micros() as u64)
.unwrap_or(0);
HlcTimestamp {
physical_micros: micros,
logical: 0,
node_tiebreaker: 0,
}
}
fn new_data_command_id() -> Result<[u8; 16], RuntimeError> {
let mut id = [0u8; 16];
getrandom::getrandom(&mut id)
.map_err(|error| RuntimeError::InvalidRequest(format!("CSPRNG failed: {error}")))?;
Ok(id)
}
struct RuntimeMetaPlane {
meta: Arc<MetaGroup<TcpTransport>>,
handle: tokio::runtime::Handle,
control: ExecutionControl,
}
impl RuntimeMetaPlane {
fn propose(&self, command: MetaCommand) -> Result<(), MetaRejectionReason> {
let meta = self.meta.clone();
let control = self.control.clone();
self.handle
.block_on(async move {
meta.propose(crate::meta::new_command_id()?, command, &control)
.await
})
.map(|_| ())
.map_err(|error| match error {
MetaError::Rejected(reason) => reason,
MetaError::Consensus(ConsensusError::NotLeader { leader }) => {
MetaRejectionReason::NotLeader { leader }
}
other => MetaRejectionReason::ProposalFailed {
reason: other.to_string(),
},
})
}
}
impl TabletMetaPlane for RuntimeMetaPlane {
fn set_tablet(&mut self, descriptor: &TabletDescriptor) -> Result<(), MetaRejectionReason> {
self.propose(MetaCommand::SetTabletDescriptor {
descriptor: descriptor.clone(),
})
}
fn tablet(&self, tablet_id: TabletId) -> Option<TabletDescriptor> {
self.meta.state().tablet(tablet_id).cloned()
}
fn remove_tablet(
&mut self,
tablet_id: TabletId,
generation: u64,
) -> Result<(), MetaRejectionReason> {
self.propose(MetaCommand::RemoveTabletDescriptor {
tablet_id,
generation,
})
}
fn publish_split(&mut self, command: &SplitPublishCommand) -> Result<(), MetaRejectionReason> {
self.propose(MetaCommand::PublishSplit {
command: command.clone(),
})
}
}
impl MergeMetaPlane for RuntimeMetaPlane {
fn publish_merge(&mut self, command: &MergePublishCommand) -> Result<(), MetaRejectionReason> {
self.propose(MetaCommand::PublishMerge {
command: command.clone(),
})
}
}
struct RuntimeKeyspace {
sink: Option<Arc<Mutex<EngineApplySink>>>,
pinned: Option<(HlcTimestamp, u64)>,
}
struct RuntimeSnapshotPin {
ts: HlcTimestamp,
_engine: EngineTabletPin,
}
impl SnapshotPin for RuntimeSnapshotPin {
fn pinned_at(&self) -> HlcTimestamp {
self.ts
}
}
impl TabletKeyspace for RuntimeKeyspace {
fn pin_snapshot(&mut self, ts: HlcTimestamp) -> Result<Box<dyn SnapshotPin>, TabletDataError> {
let pin = self
.sink
.as_ref()
.ok_or_else(|| TabletDataError::Keyspace("tablet engine is not open".to_owned()))?
.lock()
.map_err(|_| TabletDataError::Keyspace("engine sink lock poisoned".to_owned()))?
.pin_tablet_snapshot(ts)
.map_err(|error| TabletDataError::Keyspace(error.to_string()))?;
self.pinned = Some((ts, pin.epoch()));
Ok(Box::new(RuntimeSnapshotPin { ts, _engine: pin }))
}
fn snapshot_at(
&self,
ts: HlcTimestamp,
) -> Result<crate::split::RecordStream<'_>, TabletDataError> {
let (_, epoch) = self
.pinned
.filter(|(pinned, _)| *pinned == ts)
.ok_or_else(|| {
TabletDataError::Keyspace(format!("tablet snapshot {ts:?} is not pinned"))
})?;
let rows = self
.sink
.as_ref()
.ok_or_else(|| TabletDataError::Keyspace("tablet engine is not open".to_owned()))?
.lock()
.map_err(|_| TabletDataError::Keyspace("engine sink lock poisoned".to_owned()))?
.tablet_rows_at_epoch(epoch)
.map_err(|error| TabletDataError::Keyspace(error.to_string()))?;
Ok(Box::new(
rows.into_iter()
.map(|(key, value)| (Key::from_bytes(key), value)),
))
}
fn deltas_after(
&self,
ts: HlcTimestamp,
) -> Result<crate::split::MutationStream<'_>, TabletDataError> {
let (_, epoch) = self
.pinned
.filter(|(pinned, _)| *pinned == ts)
.ok_or_else(|| {
TabletDataError::Keyspace(format!("tablet snapshot {ts:?} is not pinned"))
})?;
let deltas = self
.sink
.as_ref()
.ok_or_else(|| TabletDataError::Keyspace("tablet engine is not open".to_owned()))?
.lock()
.map_err(|_| TabletDataError::Keyspace("engine sink lock poisoned".to_owned()))?
.tablet_deltas_after_epoch(epoch)
.map_err(|error| TabletDataError::Keyspace(error.to_string()))?;
Ok(Box::new(deltas.into_iter().map(
|mutation| match mutation {
TabletDataMutation::Upsert(key, value) => {
TabletMutation::Upsert(Key::from_bytes(key), value)
}
TabletDataMutation::Delete(key) => TabletMutation::Delete(Key::from_bytes(key)),
},
)))
}
}
struct RuntimeChildSink {
group: Arc<ConsensusGroup<TcpTransport>>,
handle: tokio::runtime::Handle,
control: ExecutionControl,
staged: Option<BTreeMap<Key, Vec<u8>>>,
}
impl RuntimeChildSink {
fn propose_data(&self, command: TabletDataCommand) -> Result<(), TabletDataError> {
let envelope = CommandEnvelope::new(
COMMAND_TYPE_TABLET_DATA,
new_data_command_id().map_err(|error| TabletDataError::Sink(error.to_string()))?,
TabletDataCommandRecord::new(command).encode(),
);
self.handle
.block_on(
self.group
.propose(CommandKind::Catalog, envelope, &self.control),
)
.map(|_| ())
.map_err(|error| TabletDataError::Sink(error.to_string()))
}
}
impl ChildStateSink for RuntimeChildSink {
fn begin_build(&mut self) -> Result<(), TabletDataError> {
self.staged = Some(BTreeMap::new());
Ok(())
}
fn stage(&mut self, key: &Key, value: &[u8]) -> Result<(), TabletDataError> {
let staged = self.staged.as_mut().ok_or(TabletDataError::NoStagedBuild)?;
staged.insert(key.clone(), value.to_vec());
Ok(())
}
fn install_staged(&mut self) -> Result<(), TabletDataError> {
let rows = self.staged.take().ok_or(TabletDataError::NoStagedBuild)?;
self.propose_data(TabletDataCommand::Replace {
rows: rows
.into_iter()
.map(|(key, value)| (key.into_bytes(), value))
.collect(),
})
}
fn apply_delta(&mut self, mutation: &TabletMutation) -> Result<(), TabletDataError> {
match mutation {
TabletMutation::Upsert(key, value) => self.propose_data(TabletDataCommand::Upsert {
entries: vec![(key.as_bytes().to_vec(), value.clone())],
}),
TabletMutation::Delete(key) => self.propose_data(TabletDataCommand::Delete {
keys: vec![key.as_bytes().to_vec()],
}),
}
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct HostedSyncReport {
pub created: Vec<TabletId>,
pub refreshed: Vec<TabletId>,
pub torn_down: Vec<TabletId>,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn layout_scan_tolerates_a_node_without_tablets() {
let tmp = tempfile::tempdir().unwrap();
assert_eq!(scan_tablet_layouts(tmp.path()).unwrap(), Vec::new());
std::fs::create_dir_all(tmp.path().join(TABLETS_DIR)).unwrap();
assert_eq!(scan_tablet_layouts(tmp.path()).unwrap(), Vec::new());
}
#[test]
fn layout_scan_fails_closed_on_garbage_directories() {
let tmp = tempfile::tempdir().unwrap();
let garbage = tmp.path().join(TABLETS_DIR).join("not-a-tablet");
std::fs::create_dir_all(&garbage).unwrap();
assert!(matches!(
scan_tablet_layouts(tmp.path()),
Err(RuntimeError::InvalidRequest(_))
));
std::fs::remove_dir_all(&garbage).unwrap();
let tablet_id = TabletId::new_random();
std::fs::create_dir_all(tmp.path().join(TABLETS_DIR).join(tablet_id.to_hex())).unwrap();
assert!(matches!(
scan_tablet_layouts(tmp.path()),
Err(RuntimeError::Tablet(TabletError::MissingMetadata(_)))
));
}
fn ts(micros: u64) -> HlcTimestamp {
HlcTimestamp {
physical_micros: micros,
logical: 0,
node_tiebreaker: 0,
}
}
fn pos(index: u64) -> LogPosition {
LogPosition { term: 1, index }
}
fn upsert(key: &[u8], value: &[u8]) -> TabletDataCommand {
TabletDataCommand::Upsert {
entries: vec![(key.to_vec(), value.to_vec())],
}
}
#[test]
fn tablet_ledger_applies_versions_and_partitions_the_timeline() {
let tmp = tempfile::tempdir().unwrap();
let mut ledger = TabletLedger::open(tmp.path()).unwrap();
ledger
.apply(&upsert(b"a", b"a@1"), ts(100), pos(1))
.unwrap();
ledger
.apply(&upsert(b"b", b"b@1"), ts(100), pos(2))
.unwrap();
ledger.pin(ts(100));
ledger
.apply(&upsert(b"a", b"a@2"), ts(200), pos(3))
.unwrap();
assert_eq!(
ledger.rows_at(ts(100)),
BTreeMap::from([
(Key::from_bytes(b"a".to_vec()), b"a@1".to_vec()),
(Key::from_bytes(b"b".to_vec()), b"b@1".to_vec()),
])
);
assert_eq!(
ledger.current_rows().get(&Key::from_bytes(b"a".to_vec())),
Some(&b"a@2".to_vec())
);
assert_eq!(
ledger.deltas_after(ts(100)),
vec![(Key::from_bytes(b"a".to_vec()), b"a@2".to_vec())]
);
assert!(ledger.deltas_after(ts(200)).is_empty());
ledger.unpin(ts(100));
ledger
.apply(&upsert(b"z", b"z@1"), ts(300), pos(3))
.unwrap();
assert!(!ledger
.current_rows()
.contains_key(&Key::from_bytes(b"z".to_vec())));
let position = ledger.applied_position();
drop(ledger);
let reopened = TabletLedger::open(tmp.path()).unwrap();
assert_eq!(reopened.applied_position(), position);
assert_eq!(
reopened.current_rows().get(&Key::from_bytes(b"a".to_vec())),
Some(&b"a@2".to_vec())
);
std::fs::write(
tmp.path()
.join("raft")
.join("state")
.join(TABLET_LEDGER_FILENAME),
b"junk",
)
.unwrap();
assert!(TabletLedger::open(tmp.path()).is_err());
}
#[test]
fn tablet_ledger_compacts_against_the_oldest_pin() {
let tmp = tempfile::tempdir().unwrap();
let mut ledger = TabletLedger::open(tmp.path()).unwrap();
ledger
.apply(&upsert(b"a", b"a@1"), ts(100), pos(1))
.unwrap();
ledger.pin(ts(150));
ledger
.apply(&upsert(b"a", b"a@2"), ts(200), pos(2))
.unwrap();
ledger
.apply(&upsert(b"a", b"a@3"), ts(300), pos(3))
.unwrap();
assert_eq!(
ledger.rows_at(ts(150)).get(&Key::from_bytes(b"a".to_vec())),
Some(&b"a@1".to_vec())
);
assert_eq!(
ledger.deltas_after(ts(150)),
vec![
(Key::from_bytes(b"a".to_vec()), b"a@2".to_vec()),
(Key::from_bytes(b"a".to_vec()), b"a@3".to_vec()),
]
);
ledger.unpin(ts(150));
ledger
.apply(&upsert(b"a", b"a@4"), ts(400), pos(4))
.unwrap();
assert_eq!(
ledger.rows_at(ts(150)).get(&Key::from_bytes(b"a".to_vec())),
None
);
assert_eq!(
ledger.current_rows().get(&Key::from_bytes(b"a".to_vec())),
Some(&b"a@4".to_vec())
);
assert_eq!(ledger.pin_count(), 0);
}
#[test]
fn tablet_ledger_replace_is_atomic_and_refused_under_a_live_pin() {
let tmp = tempfile::tempdir().unwrap();
let mut ledger = TabletLedger::open(tmp.path()).unwrap();
ledger
.apply(&upsert(b"a", b"a@1"), ts(100), pos(1))
.unwrap();
ledger.pin(ts(150));
let replace = TabletDataCommand::Replace {
rows: vec![(b"b".to_vec(), b"b@2".to_vec())],
};
assert!(ledger.apply(&replace, ts(200), pos(2)).is_err());
ledger.unpin(ts(150));
ledger.apply(&replace, ts(200), pos(2)).unwrap();
assert_eq!(
ledger.current_rows(),
BTreeMap::from([(Key::from_bytes(b"b".to_vec()), b"b@2".to_vec())])
);
let bytes = ledger.snapshot_bytes().unwrap();
let follower_dir = tempfile::tempdir().unwrap();
let mut follower = TabletLedger::open(follower_dir.path()).unwrap();
follower.install_bytes(&bytes).unwrap();
assert_eq!(follower.current_rows(), ledger.current_rows());
assert!(follower.install_bytes(b"junk").is_err());
}
#[test]
fn group_snapshot_frame_round_trips_and_fails_closed() {
let engine = b"engine-half".to_vec();
let ledger = b"ledger-half".to_vec();
let frame = encode_group_snapshot(&engine, &ledger);
let (engine_back, ledger_back) = decode_group_snapshot(&frame).unwrap();
assert_eq!(engine_back, engine);
assert_eq!(ledger_back, ledger);
let (empty_engine, empty_ledger) =
decode_group_snapshot(&encode_group_snapshot(&[], &[])).unwrap();
assert!(empty_engine.is_empty() && empty_ledger.is_empty());
assert!(decode_group_snapshot(&frame[..6]).is_err());
assert!(decode_group_snapshot(&frame[..12 + engine.len() - 1]).is_err());
let mut future = frame.clone();
future[..4].copy_from_slice(&99_u32.to_le_bytes());
assert!(decode_group_snapshot(&future).is_err());
}
#[test]
fn tablet_data_command_record_round_trips_and_fails_closed() {
let record = TabletDataCommandRecord::new(upsert(b"k", b"v"));
let decoded = TabletDataCommandRecord::decode(&record.encode()).unwrap();
assert_eq!(decoded, record);
assert!(TabletDataCommandRecord::decode(b"not json").is_err());
let mut value: serde_json::Value = serde_json::from_slice(&record.encode()).unwrap();
value["format_version"] = serde_json::json!(99);
assert!(TabletDataCommandRecord::decode(&serde_json::to_vec(&value).unwrap()).is_err());
}
}