use std::collections::{BTreeMap, BTreeSet};
use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, SystemTime};
use async_trait::async_trait;
use tokio::sync::{mpsc, oneshot, watch};
use super::bootstrap::{
inspect, BootDecision, BootstrapEffect, Integrity, QuarantinedNode, RecoveredState,
RejoinMessage,
};
use super::driver::{
run_apply_worker, spawn_ticker, BarrierOutcome, DriverConfig, DriverEvent, DriverExit,
DurableState, FileStore, ProposeOutcome, Transport, WireMsg, YrpDriver, YrpStatus,
};
use super::engine_sink::{AppliedOutcome, EngineApplySink, OutcomeStore};
use super::op::{claim_key_for_op, YrpOp};
use super::replica::{LogEntry, Payload, Role};
use super::transport::HttpTransport;
use super::types::{ClusterId, NodeId};
use crate::commit::{
Applier, CommitError, CommitOptions, CommitReceipt, CommittedEntry, MemoryMutation,
MutationCommitter, OpId, TenantId,
};
#[derive(Debug, Clone)]
pub struct YrpPeer {
pub node_id: u64,
pub addr: String,
pub witness: bool,
}
#[derive(Debug, Clone)]
pub struct YrpRuntimeConfig {
pub node_id: u64,
pub cluster_id: u64,
pub peers: Vec<YrpPeer>,
pub data_dir: PathBuf,
pub cluster_secret: Option<String>,
pub tick_ms: u64,
pub election_ticks: (u32, u32),
pub heartbeat_ticks: u32,
pub compact_after_entries: u64,
pub leader_retain_entries: u64,
}
const PROPOSE_TIMEOUT: Duration = Duration::from_secs(15);
const OUTCOME_POLL: Duration = Duration::from_millis(10);
const REJOIN_RETRY: Duration = Duration::from_secs(2);
#[derive(Debug)]
pub enum YrpProposeError {
NotLeader {
leader_id: Option<u64>,
leader_addr: Option<String>,
},
Timeout,
Unavailable(String),
}
#[derive(Debug)]
pub enum ControlWriteError {
AlreadyExists,
Diverged(String),
Internal(String),
Propose(YrpProposeError),
}
pub struct YrpHandle {
pub node_id: NodeId,
pub cluster_id: u64,
owner_tx: mpsc::UnboundedSender<DriverEvent>,
pub outcomes: Arc<OutcomeStore>,
pub status: watch::Receiver<YrpStatus>,
peer_http: BTreeMap<u64, String>,
pub cluster_secret: Option<String>,
local: Arc<dyn MutationCommitter>,
control: Arc<parking_lot::Mutex<crate::control::ControlDb>>,
control_propose_lock: tokio::sync::Mutex<()>,
quarantine: std::sync::RwLock<Option<Vec<String>>>,
}
impl YrpHandle {
pub fn quarantine_reasons(&self) -> Option<Vec<String>> {
self.quarantine.read().expect("quarantine lock").clone()
}
pub fn peer_urls(&self) -> Vec<(u64, String)> {
self.peer_http
.iter()
.map(|(id, url)| (*id, url.clone()))
.collect()
}
fn set_quarantine(&self, reasons: Option<Vec<String>>) {
*self.quarantine.write().expect("quarantine lock") = reasons;
}
pub fn deliver(&self, from: u64, msg: WireMsg) -> Result<(), String> {
super::transport::deliver(&self.owner_tx, from, msg)
}
pub async fn read_barrier(&self) -> Result<(), YrpProposeError> {
if let Some(reasons) = self.quarantine_reasons() {
return Err(YrpProposeError::Unavailable(format!(
"node quarantined: {reasons:?}"
)));
}
let (tx, rx) = oneshot::channel();
self.owner_tx
.send(DriverEvent::ReadBarrier { reply: tx })
.map_err(|_| YrpProposeError::Unavailable("YRP driver not running".into()))?;
let out = tokio::time::timeout(PROPOSE_TIMEOUT, rx)
.await
.map_err(|_| YrpProposeError::Timeout)?
.map_err(|_| YrpProposeError::Unavailable("YRP driver dropped reply".into()))?;
match out {
BarrierOutcome::Ok => Ok(()),
BarrierOutcome::Retry => {
let (leader_id, leader_addr) = self.leader_hint();
Err(YrpProposeError::NotLeader {
leader_id,
leader_addr,
})
}
}
}
pub async fn serve_backfill(
&self,
from_index: u64,
to_index: u64,
) -> Result<Vec<(u64, LogEntry)>, String> {
if self.status.borrow().engine_incomplete() {
return Err("node is engine-incomplete; cannot source backfill".into());
}
let outcomes = self
.outcomes
.outcomes_in_range(from_index, to_index)
.map_err(|e| format!("outcome range: {e}"))?;
let expected = (to_index.saturating_sub(from_index)) as usize;
if outcomes.len() != expected {
return Err(format!(
"backfill range ({from_index},{to_index}] not fully retained: {} of {expected} rows",
outcomes.len()
));
}
let mut rows = Vec::with_capacity(outcomes.len());
for o in outcomes {
let tenant = TenantId::new(o.tenant_id);
let entry = self
.local
.read_range(tenant, o.tenant_log_index, 1)
.await
.map_err(|e| format!("commit-log read: {e}"))?
.into_iter()
.next()
.ok_or_else(|| {
format!(
"commit-log row missing for tenant {tenant} idx {}",
o.tenant_log_index
)
})?;
let op_id = OpId::from_uuid(
o.op_id
.parse()
.map_err(|e| format!("bad op_id in outcome: {e}"))?,
);
let op = YrpOp {
tenant_id: tenant,
op_id,
mutation: entry.mutation,
idempotency_key: o.key_str.clone(),
};
let log_entry = LogEntry {
term: super::types::Term(o.term),
payload: Payload::Op(op.encode()?),
key: o.key_hash,
activate: None,
};
rows.push((o.yrp_index, log_entry));
}
Ok(rows)
}
pub fn shutdown(&self) {
let _ = self.owner_tx.send(DriverEvent::Shutdown);
}
pub fn is_stopped(&self) -> bool {
self.owner_tx.is_closed()
}
pub fn leader_hint(&self) -> (Option<u64>, Option<String>) {
let leader = self.status.borrow().leader.map(|n| n.0);
let addr = leader.and_then(|id| self.peer_http.get(&id).cloned());
(leader, addr)
}
pub fn is_leader(&self) -> bool {
let s = *self.status.borrow();
self.quarantine_reasons().is_none() && s.role == Role::Leader && !s.engine_incomplete()
}
pub fn engine_incomplete(&self) -> bool {
self.status.borrow().engine_incomplete()
}
pub async fn propose_and_wait(
&self,
key: u64,
op: &YrpOp,
) -> Result<AppliedOutcome, YrpProposeError> {
if let Some(reasons) = self.quarantine_reasons() {
return Err(YrpProposeError::Unavailable(format!(
"node quarantined: {reasons:?}"
)));
}
let bytes = op.encode().map_err(YrpProposeError::Unavailable)?;
let (tx, rx) = oneshot::channel();
self.owner_tx
.send(DriverEvent::Propose {
key,
payload: Payload::Op(bytes),
reply: tx,
})
.map_err(|_| YrpProposeError::Unavailable("YRP driver not running".into()))?;
let outcome = tokio::time::timeout(PROPOSE_TIMEOUT, rx)
.await
.map_err(|_| YrpProposeError::Timeout)?
.map_err(|_| YrpProposeError::Unavailable("YRP driver dropped reply".into()))?;
let index = match outcome {
ProposeOutcome::Applied { index } | ProposeOutcome::Duplicate { index } => index,
ProposeOutcome::Retry => {
let (leader_id, leader_addr) = self.leader_hint();
return Err(YrpProposeError::NotLeader {
leader_id,
leader_addr,
});
}
};
self.wait_outcome(index).await
}
pub async fn propose_control(
&self,
actor: &str,
op: &super::control_op::ControlOp,
) -> Result<u64, YrpProposeError> {
if let Some(reasons) = self.quarantine_reasons() {
return Err(YrpProposeError::Unavailable(format!(
"node quarantined: {reasons:?}"
)));
}
let env = super::control_op::ControlEnvelope::new(actor, op.clone());
let bytes = env.encode().map_err(YrpProposeError::Unavailable)?;
let (tx, rx) = oneshot::channel();
self.owner_tx
.send(DriverEvent::Propose {
key: env.claim_key(),
payload: Payload::Control(bytes),
reply: tx,
})
.map_err(|_| YrpProposeError::Unavailable("YRP driver not running".into()))?;
let outcome = tokio::time::timeout(PROPOSE_TIMEOUT, rx)
.await
.map_err(|_| YrpProposeError::Timeout)?
.map_err(|_| YrpProposeError::Unavailable("YRP driver dropped reply".into()))?;
let index = match outcome {
ProposeOutcome::Applied { index } | ProposeOutcome::Duplicate { index } => index,
ProposeOutcome::Retry => {
let (leader_id, leader_addr) = self.leader_hint();
return Err(YrpProposeError::NotLeader {
leader_id,
leader_addr,
});
}
};
let deadline = tokio::time::Instant::now() + PROPOSE_TIMEOUT;
while self.outcomes.applied() < index {
if tokio::time::Instant::now() >= deadline {
return Err(YrpProposeError::Timeout);
}
tokio::time::sleep(OUTCOME_POLL).await;
}
Ok(index)
}
pub async fn create_database_replicated(
&self,
actor: &str,
name: &str,
path: &str,
config: &str,
created_at: String,
) -> Result<i64, ControlWriteError> {
let _guard = self.control_propose_lock.lock().await;
let db_id = {
let db = self.control.lock();
if db
.database_exists(name)
.map_err(|e| ControlWriteError::Internal(format!("name check: {e}")))?
{
return Err(ControlWriteError::AlreadyExists);
}
db.next_database_id()
.map_err(|e| ControlWriteError::Internal(format!("allocate db id: {e}")))?
};
let op = super::control_op::ControlOp::CreateDatabase {
db_id,
name: name.to_string(),
path: path.to_string(),
config: config.to_string(),
created_at,
};
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)?;
match self
.control
.lock()
.get_database(name)
.map_err(|e| ControlWriteError::Internal(format!("read-back: {e}")))?
{
Some(rec) => Ok(rec.id),
None => Err(ControlWriteError::Diverged(format!(
"CreateDatabase({name}) committed but no row after apply"
))),
}
}
pub async fn create_token_replicated(
&self,
actor: &str,
db_id: i64,
token_hash: String,
label: String,
created_at: String,
) -> Result<(), ControlWriteError> {
let op = super::control_op::ControlOp::CreateToken {
db_id,
token_hash: token_hash.clone(),
label,
created_at,
};
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)?;
let resolved = self
.control
.lock()
.validate_token(&token_hash)
.map_err(|e| ControlWriteError::Internal(format!("verify token: {e}")))?;
if resolved == Some(db_id) {
Ok(())
} else {
Err(ControlWriteError::Diverged(
"CreateToken committed but token does not resolve after apply".into(),
))
}
}
pub async fn revoke_token_replicated(
&self,
actor: &str,
token_hash: String,
revoked_at: String,
) -> Result<(), ControlWriteError> {
let op = super::control_op::ControlOp::RevokeToken {
token_hash: token_hash.clone(),
revoked_at,
};
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)?;
let resolved = self
.control
.lock()
.validate_token(&token_hash)
.map_err(|e| ControlWriteError::Internal(format!("verify revoke: {e}")))?;
if resolved.is_none() {
Ok(())
} else {
Err(ControlWriteError::Diverged(
"RevokeToken committed but token still resolves after apply".into(),
))
}
}
pub async fn create_user_replicated(
&self,
actor: &str,
username: &str,
password_hash: String,
role: String,
created_at: String,
) -> Result<(), ControlWriteError> {
let _guard = self.control_propose_lock.lock().await;
if self
.control
.lock()
.get_admin_user(username)
.map_err(|e| ControlWriteError::Internal(format!("user check: {e}")))?
.is_some()
{
return Err(ControlWriteError::AlreadyExists);
}
let op = super::control_op::ControlOp::CreateUser {
username: username.to_string(),
password_hash,
role,
created_at,
};
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)?;
match self.control.lock().get_admin_user(username) {
Ok(Some(_)) => Ok(()),
_ => Err(ControlWriteError::Diverged(
"CreateUser committed but user absent after apply".into(),
)),
}
}
pub async fn set_user_role_replicated(
&self,
actor: &str,
username: &str,
role: String,
) -> Result<(), ControlWriteError> {
let op = super::control_op::ControlOp::SetUserRole {
username: username.to_string(),
role,
};
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)
.map(|_| ())
}
pub async fn set_user_password_replicated(
&self,
actor: &str,
username: &str,
password_hash: String,
) -> Result<(), ControlWriteError> {
let op = super::control_op::ControlOp::SetUserPassword {
username: username.to_string(),
password_hash,
};
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)
.map(|_| ())
}
pub async fn disable_user_replicated(
&self,
actor: &str,
username: &str,
disabled_at: String,
) -> Result<(), ControlWriteError> {
let op = super::control_op::ControlOp::DisableUser {
username: username.to_string(),
disabled_at,
};
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)
.map(|_| ())
}
pub async fn set_admin_session_key_replicated(
&self,
actor: &str,
kid: String,
value: String,
) -> Result<(), ControlWriteError> {
let op = super::control_op::ControlOp::SetAdminSessionKey { kid, value };
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)
.map(|_| ())
}
pub async fn mount_pack_replicated(
&self,
actor: &str,
database_id: i64,
pack_digest: String,
pack_name: String,
mounted_at: String,
) -> Result<(), ControlWriteError> {
let op = super::control_op::ControlOp::MountPack {
database_id,
pack_digest,
pack_name,
mounted_at,
};
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)
.map(|_| ())
}
pub async fn unmount_pack_replicated(
&self,
actor: &str,
database_id: i64,
pack_digest: String,
unmounted_at: String,
) -> Result<(), ControlWriteError> {
let op = super::control_op::ControlOp::UnmountPack {
database_id,
pack_digest,
unmounted_at,
};
self.propose_control(actor, &op)
.await
.map_err(ControlWriteError::Propose)
.map(|_| ())
}
async fn wait_outcome(&self, index: u64) -> Result<AppliedOutcome, YrpProposeError> {
let deadline = tokio::time::Instant::now() + PROPOSE_TIMEOUT;
while self.outcomes.applied() < index {
if tokio::time::Instant::now() >= deadline {
return Err(YrpProposeError::Timeout);
}
tokio::time::sleep(OUTCOME_POLL).await;
}
self.outcomes
.lookup_by_index(index)
.map_err(YrpProposeError::Unavailable)?
.ok_or_else(|| {
YrpProposeError::Unavailable(format!("applied index {index} has no outcome record"))
})
}
}
pub fn spawn(
cfg: YrpRuntimeConfig,
local: Arc<dyn MutationCommitter>,
applier: Arc<dyn Applier>,
control: Arc<parking_lot::Mutex<crate::control::ControlDb>>,
) -> Result<Arc<YrpHandle>, String> {
if cfg.cluster_id == 0 {
return Err("[yrp] cluster_id must be non-zero".into());
}
if cfg.node_id == 0 {
return Err("[cluster] node_id must be non-zero in yrp mode".into());
}
if !cfg.peers.iter().any(|p| p.node_id == cfg.node_id) {
return Err("[yrp] peers must include this node's node_id".into());
}
if cfg.compact_after_entries > 0 {
tracing::warn!(
compact_after = cfg.compact_after_entries,
"[yrp] log compaction ENABLED: beyond-GC stragglers rejoin without \
engine backfill for the compacted range until Phase C \
(engine-checkpoint transfer). Not recommended in production."
);
tracing::error!(
compact_after = cfg.compact_after_entries,
"[yrp] compaction + RFC 029 control-plane replication is NOT safe: \
control ops in a compacted range are lost to rejoiners. Set \
compact_after_entries = 0 until RFC 029 increment 2 (control \
state in the snapshot) ships."
);
}
let me = NodeId(cfg.node_id);
let cluster = ClusterId(cfg.cluster_id);
let state_path = cfg.data_dir.join("yrp.state");
let outcomes = Arc::new(OutcomeStore::open(cfg.data_dir.join("yrp_apply.sqlite"))?);
let voters: BTreeSet<NodeId> = cfg.peers.iter().map(|p| NodeId(p.node_id)).collect();
let witnesses: BTreeSet<NodeId> = cfg
.peers
.iter()
.filter(|p| p.witness)
.map(|p| NodeId(p.node_id))
.collect();
let peer_http: BTreeMap<u64, String> = cfg
.peers
.iter()
.map(|p| (p.node_id, p.addr.clone()))
.collect();
let peer_urls: BTreeMap<NodeId, String> = cfg
.peers
.iter()
.filter(|p| p.node_id != cfg.node_id)
.map(|p| (NodeId(p.node_id), p.addr.clone()))
.collect();
let data_peers: Vec<NodeId> = cfg
.peers
.iter()
.filter(|p| p.node_id != cfg.node_id && !p.witness)
.map(|p| NodeId(p.node_id))
.collect();
let transport = Arc::new(HttpTransport::new(
me,
peer_urls,
cfg.cluster_secret.clone(),
));
let store = FileStore::new(state_path.clone());
let (restored, recovered) = match store.load() {
Ok(Some(d)) => {
let rec = RecoveredState {
cluster_id: Some(d.cluster_id),
hard: Some(d.hard),
log: Some(d.log.clone()),
active: d.active,
commit_marker: outcomes.applied().saturating_sub(d.base.index),
integrity: Integrity {
hard_state_verified: true,
log_verified: true,
},
};
(Some(d), rec)
}
Ok(None) => {
let rec = RecoveredState {
cluster_id: None,
hard: None,
log: None,
active: 0,
commit_marker: outcomes.applied(),
integrity: Integrity {
hard_state_verified: true,
log_verified: true,
},
};
(None, rec)
}
Err(e) => {
tracing::error!(error = %e, "yrp.state unreadable — boot inspection will quarantine");
let rec = RecoveredState {
cluster_id: None,
hard: None,
log: None,
active: 0,
commit_marker: outcomes.applied(),
integrity: Integrity {
hard_state_verified: false,
log_verified: false,
},
};
(None, rec)
}
};
let (owner_tx, owner_rx) = mpsc::unbounded_channel();
let (apply_tx, apply_rx) = mpsc::unbounded_channel();
let (status_tx, status_rx) = watch::channel(YrpStatus::default());
let handle = Arc::new(YrpHandle {
node_id: me,
cluster_id: cfg.cluster_id,
owner_tx: owner_tx.clone(),
outcomes: outcomes.clone(),
status: status_rx,
peer_http,
cluster_secret: cfg.cluster_secret.clone(),
local: local.clone(),
control: control.clone(),
control_propose_lock: tokio::sync::Mutex::new(()),
quarantine: std::sync::RwLock::new(None),
});
tokio::spawn(run_backfill_task(handle.clone(), owner_tx.clone()));
let control_sink = Arc::new(super::control_op::ControlApplySink::new(control));
let sink = EngineApplySink::new(local, applier, outcomes.clone()).with_control(control_sink);
tokio::spawn(run_apply_worker(Box::new(sink), apply_rx, owner_tx.clone()));
spawn_ticker(owner_tx.clone(), Duration::from_millis(cfg.tick_ms.max(1)));
let driver_cfg = move || DriverConfig {
id: me,
cluster_id: cluster,
voters: voters.clone(),
witnesses: witnesses.clone(),
supported: u32::MAX,
election_ticks: cfg.election_ticks,
heartbeat_ticks: cfg.heartbeat_ticks,
compact_after: (cfg.compact_after_entries > 0).then_some(cfg.compact_after_entries),
leader_retain: cfg.leader_retain_entries,
};
match inspect(cluster, u32::MAX, &recovered) {
BootDecision::Healthy { .. } => {
tracing::info!(
node = cfg.node_id,
cluster = cfg.cluster_id,
"YRP boot: healthy"
);
let mut driver = YrpDriver::new(
driver_cfg(),
restored,
store,
Box::new(SharedTransport(transport)),
apply_tx,
outcomes.applied(),
);
driver.set_status_tx(status_tx);
spawn_driver(driver, owner_rx, handle.clone());
}
BootDecision::Quarantine { reasons, term_hint } => {
tracing::error!(
?reasons,
"YRP boot: QUARANTINED (fail closed, serving diagnostics)"
);
handle.set_quarantine(Some(reasons.iter().map(|r| format!("{r:?}")).collect()));
let node = QuarantinedNode::new(me, cluster, reasons, term_hint);
let ctx = QuarantineCtx {
node,
store,
state_path,
transport,
apply_tx,
status_tx,
data_peers,
outcomes,
driver_cfg: driver_cfg(),
handle: handle.clone(),
};
tokio::spawn(run_quarantined(ctx, owner_rx));
}
}
Ok(handle)
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct BackfillRequest {
pub cluster_id: u64,
pub from_index: u64,
pub to_index: u64,
}
const BACKFILL_BATCH: u64 = 256;
async fn run_backfill_task(handle: Arc<YrpHandle>, owner: mpsc::UnboundedSender<DriverEvent>) {
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(30))
.build()
.expect("reqwest client");
let mut status = handle.status.clone();
loop {
let (from, to, leader_id) = {
let s = *status.borrow();
(s.applied, s.backfill_target, s.leader.map(|n| n.0))
};
if to <= from {
if status.changed().await.is_err() {
return;
}
continue;
}
let Some(leader_id) = leader_id.filter(|id| *id != handle.node_id.0) else {
if status.changed().await.is_err() {
return;
}
continue;
};
let Some(base_url) = handle.peer_http.get(&leader_id).cloned() else {
tokio::time::sleep(Duration::from_millis(200)).await;
continue;
};
let batch_to = (from + BACKFILL_BATCH).min(to);
let req = BackfillRequest {
cluster_id: handle.cluster_id,
from_index: from,
to_index: batch_to,
};
let url = format!("{}/v1/yrp/backfill", base_url.trim_end_matches('/'));
let mut http = client.post(&url).json(&req);
if let Some(s) = &handle.cluster_secret {
http = http.bearer_auth(s);
}
match http.send().await {
Ok(resp) if resp.status().is_success() => match resp.bytes().await {
Ok(bytes) => match bincode::deserialize::<Vec<(u64, LogEntry)>>(&bytes) {
Ok(rows) => {
for (index, entry) in rows {
let _ = owner.send(DriverEvent::Backfilled { index, entry });
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
Err(e) => {
tracing::warn!(error = %e, "backfill decode failed; retrying");
tokio::time::sleep(Duration::from_millis(200)).await;
}
},
Err(e) => {
tracing::warn!(error = %e, "backfill body read failed; retrying");
tokio::time::sleep(Duration::from_millis(200)).await;
}
},
other => {
tracing::debug!(?other, from, batch_to, "backfill pull failed; retrying");
tokio::time::sleep(Duration::from_millis(300)).await;
}
}
}
}
struct SharedTransport(Arc<HttpTransport>);
impl Transport for SharedTransport {
fn send(&self, to: NodeId, msg: WireMsg) {
self.0.send(to, msg)
}
}
fn spawn_driver(
driver: YrpDriver,
owner_rx: mpsc::UnboundedReceiver<DriverEvent>,
handle: Arc<YrpHandle>,
) {
tokio::spawn(async move {
let exit = driver.run(owner_rx).await;
match exit {
DriverExit::Shutdown => {
tracing::info!("YRP driver shut down");
}
other => {
tracing::error!(
?other,
"YRP driver FAILED — node degraded to quarantine posture"
);
handle.set_quarantine(Some(vec![format!("driver exit: {other:?}")]));
}
}
});
}
struct QuarantineCtx {
node: QuarantinedNode,
store: FileStore,
state_path: PathBuf,
transport: Arc<HttpTransport>,
apply_tx: mpsc::UnboundedSender<(u64, super::replica::LogEntry)>,
status_tx: watch::Sender<YrpStatus>,
data_peers: Vec<NodeId>,
outcomes: Arc<OutcomeStore>,
driver_cfg: DriverConfig,
handle: Arc<YrpHandle>,
}
async fn run_quarantined(mut ctx: QuarantineCtx, mut rx: mpsc::UnboundedReceiver<DriverEvent>) {
let mut retry = tokio::time::interval(REJOIN_RETRY);
retry.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut peer_rr = 0usize;
loop {
tokio::select! {
_ = retry.tick() => {
if ctx.data_peers.is_empty() {
continue;
}
let target = ctx.data_peers[peer_rr % ctx.data_peers.len()];
peer_rr += 1;
for eff in ctx.node.tick_rejoin(target) {
run_bootstrap_effect(&mut ctx, eff);
}
}
ev = rx.recv() => {
let Some(ev) = ev else { return };
match ev {
DriverEvent::Shutdown => return,
DriverEvent::Inbound { from, msg: WireMsg::Rejoin(grant @ RejoinMessage::Grant { .. }) } => {
for eff in ctx.node.on_grant(from, grant.clone()) {
if let BootstrapEffect::AdoptSnapshot { cluster_id, hard, base, log, claims, active } = eff {
let adopted = DurableState { cluster_id, hard, base, log, claims, active };
if let Err(e) = ctx.store.persist(&adopted) {
tracing::error!(error = %e, "adopt-snapshot persist failed; staying quarantined");
continue;
}
tracing::info!(?from, "YRP rejoin: snapshot adopted — resuming as follower");
ctx.handle.set_quarantine(None);
let mut driver = YrpDriver::new(
ctx.driver_cfg,
Some(adopted),
ctx.store,
Box::new(SharedTransport(ctx.transport)),
ctx.apply_tx,
ctx.outcomes.applied(),
);
driver.set_status_tx(ctx.status_tx);
spawn_driver(driver, rx, ctx.handle);
return;
}
}
}
_ => {}
}
}
}
}
}
fn run_bootstrap_effect(ctx: &mut QuarantineCtx, eff: BootstrapEffect) {
match eff {
BootstrapEffect::PreserveOldState => {
let ts = SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let dst = ctx.state_path.with_extension(format!("preserved-{ts}"));
match std::fs::copy(&ctx.state_path, &dst) {
Ok(_) => tracing::warn!(dst = %dst.display(), "quarantine: old state preserved"),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => tracing::error!(error = %e, "quarantine: preserve-old-state failed"),
}
}
BootstrapEffect::Alarm { reasons } => {
tracing::error!(
?reasons,
"YRP QUARANTINE ALARM: corruption evidence — operator attention required"
);
}
BootstrapEffect::Send { to, msg } => {
ctx.transport.send(to, WireMsg::Rejoin(msg));
}
BootstrapEffect::AdoptSnapshot { .. } => {
unreachable!("adopt is handled inline by run_quarantined")
}
}
}
pub struct YrpCommitter {
handle: Arc<YrpHandle>,
local: Arc<dyn MutationCommitter>,
}
impl YrpCommitter {
pub fn new(handle: Arc<YrpHandle>, local: Arc<dyn MutationCommitter>) -> Self {
Self { handle, local }
}
}
pub fn propose_err_to_commit(e: YrpProposeError, op_id: OpId) -> CommitError {
match e {
YrpProposeError::NotLeader {
leader_id,
leader_addr,
} => CommitError::NotLeader {
leader_id,
leader_addr,
},
YrpProposeError::Timeout => CommitError::CommitTimeout { op_id },
YrpProposeError::Unavailable(m) => CommitError::StorageFailure { message: m },
}
}
#[async_trait]
impl MutationCommitter for YrpCommitter {
async fn commit(
&self,
tenant_id: TenantId,
mutation: MemoryMutation,
opts: CommitOptions,
) -> Result<CommitReceipt, CommitError> {
if let Some(expected) = opts.expected_log_index {
return Err(CommitError::StorageFailure {
message: format!("expected_log_index ({expected}) is not supported in yrp mode"),
});
}
if !mutation.is_implemented() {
return Err(CommitError::NotYetImplemented {
variant: mutation.variant_name(),
planned_rfc: mutation.planned_rfc(),
});
}
let mut op_id = opts.op_id.unwrap_or_else(OpId::new_random);
for attempt in 0..3 {
let op = YrpOp {
tenant_id,
op_id,
mutation: mutation.clone(),
idempotency_key: None,
};
let key = claim_key_for_op(tenant_id, &op_id);
let outcome = self
.handle
.propose_and_wait(key, &op)
.await
.map_err(|e| propose_err_to_commit(e, op_id))?;
if outcome.op_id != op_id.to_string() {
tracing::warn!(
attempt,
"yrp unkeyed claim digest collision detected; retrying with fresh op_id"
);
if opts.op_id.is_some() {
return Err(CommitError::StorageFailure {
message: "claim digest collision on caller-supplied op_id".into(),
});
}
op_id = OpId::new_random();
continue;
}
let applied_at = SystemTime::UNIX_EPOCH
+ Duration::from_micros(outcome.applied_at_unix_micros.max(0) as u64);
return Ok(CommitReceipt {
op_id,
tenant_id,
term: outcome.term,
log_index: outcome.tenant_log_index,
committed_at: applied_at,
applied_at: Some(applied_at),
});
}
Err(CommitError::StorageFailure {
message: "repeated claim digest collisions (unkeyed)".into(),
})
}
async fn read_range(
&self,
tenant_id: TenantId,
from_index: u64,
limit: usize,
) -> Result<Vec<CommittedEntry>, CommitError> {
self.local.read_range(tenant_id, from_index, limit).await
}
async fn high_watermark(&self, tenant_id: TenantId) -> Result<u64, CommitError> {
self.local.high_watermark(tenant_id).await
}
async fn list_active_tenants(&self) -> Result<Vec<TenantId>, CommitError> {
self.local.list_active_tenants().await
}
async fn ensure_linearizable(&self) -> Result<(), CommitError> {
self.handle
.read_barrier()
.await
.map_err(|e| propose_err_to_commit(e, OpId::new_random()))
}
}