use std::collections::BTreeSet;
use std::time::Duration;
use nodedb_cluster::calvin::SEQUENCER_GROUP_ID;
use nodedb_cluster::calvin::types::{LockKeyWire, ReleaseReason, TxnIdWire};
use nodedb_cluster::{
RaftRpc, ReleaseReservationRequest, ReleaseReservationResponse, ReserveReadRequest,
ReserveReadResponse,
};
use crate::Error;
use crate::control::server::exchange::resolve::register_peers_from_topology;
use crate::control::state::SharedState;
pub(crate) async fn submit_local_reserve_read(
state: &SharedState,
key: LockKeyWire,
vshard: u32,
owner: Option<TxnIdWire>,
timeout: Duration,
) -> crate::Result<TxnIdWire> {
let inbox = state
.reservation_inbox
.get()
.ok_or(Error::SequencerUnavailable)?;
let rx = inbox
.submit_reserve(key, vshard, owner)
.map_err(|e| Error::BadRequest {
detail: format!("reservation service rejected reserve-read: {e}"),
})?;
tokio::time::timeout(timeout, rx)
.await
.map_err(|_| Error::Internal {
detail: "timed out waiting for reservation assignment".to_owned(),
})?
.map_err(|_| Error::Internal {
detail: "reservation assignment channel closed".to_owned(),
})
}
pub(crate) async fn submit_reserve_read(
state: &SharedState,
key: LockKeyWire,
vshard: u32,
owner: Option<TxnIdWire>,
) -> crate::Result<TxnIdWire> {
let local_timeout = Duration::from_secs(state.tuning.network.default_deadline_secs);
let (Some(transport), Some(_routing)) = (
state.cluster_transport.as_ref(),
state.cluster_routing.as_ref(),
) else {
return submit_local_reserve_read(state, key, vshard, owner, local_timeout).await;
};
let status_fn = state.raft_status_fn.get().ok_or_else(|| Error::Internal {
detail: "reserve-read: raft status fn not installed (cluster not started)".to_owned(),
})?;
let leader = status_fn()
.into_iter()
.find(|g| g.group_id == SEQUENCER_GROUP_ID)
.map(|g| g.leader_id)
.unwrap_or(0);
if leader == 0 {
return Err(Error::Internal {
detail: "no sequencer leader elected yet; cannot reserve read".to_owned(),
});
}
if leader == state.node_id {
return submit_local_reserve_read(state, key, vshard, owner, local_timeout).await;
}
let mut targets = BTreeSet::new();
targets.insert(leader);
register_peers_from_topology(state, transport, &targets);
let lock_key_bytes = zerompk::to_msgpack_vec(&key).map_err(|e| Error::Serialization {
format: "msgpack".to_owned(),
detail: format!("failed to encode LockKeyWire for routed reserve-read: {e}"),
})?;
let owner_bytes = owner
.map(|o| {
zerompk::to_msgpack_vec(&o).map_err(|e| Error::Serialization {
format: "msgpack".to_owned(),
detail: format!("failed to encode TxnIdWire owner for routed reserve-read: {e}"),
})
})
.transpose()?;
let deadline_remaining_ms = state
.tuning
.network
.default_deadline_secs
.saturating_mul(1000)
.max(1);
let req = ReserveReadRequest {
lock_key_bytes,
vshard,
owner_bytes,
deadline_remaining_ms,
trace_id: [0u8; 16],
};
let read_timeout = Duration::from_millis(deadline_remaining_ms.saturating_add(2_000));
match transport
.send_rpc_with_read_timeout(leader, RaftRpc::ReserveReadRequest(req), read_timeout)
.await
{
Ok(RaftRpc::ReserveReadResponse(ReserveReadResponse {
owner_bytes: Some(b),
error: None,
})) => zerompk::from_msgpack::<TxnIdWire>(&b).map_err(|e| Error::Serialization {
format: "msgpack".to_owned(),
detail: format!("failed to decode TxnIdWire owner from reserve-read reply: {e}"),
}),
Ok(RaftRpc::ReserveReadResponse(ReserveReadResponse { error: Some(e), .. })) => {
Err(Error::Internal {
detail: format!("reserve-read failed on sequencer leader node {leader}: {e:?}"),
})
}
Ok(other) => Err(Error::Internal {
detail: format!("reserve-read: unexpected reply from node {leader}: {other:?}"),
}),
Err(e) => Err(Error::Internal {
detail: format!("reserve-read RPC to sequencer leader node {leader} failed: {e}"),
}),
}
}
pub(crate) async fn submit_local_release(
state: &SharedState,
owner: TxnIdWire,
vshard: u32,
reason: ReleaseReason,
) -> crate::Result<()> {
state
.reservation_inbox
.get()
.ok_or(Error::SequencerUnavailable)?
.submit_release(owner, vshard, reason)
.map_err(|e| Error::BadRequest {
detail: format!("reservation release rejected: {e}"),
})?;
Ok(())
}
pub(crate) async fn release_reservation(
state: &SharedState,
owner: TxnIdWire,
vshard: u32,
reason: ReleaseReason,
) -> crate::Result<()> {
let (Some(transport), Some(_routing)) = (
state.cluster_transport.as_ref(),
state.cluster_routing.as_ref(),
) else {
return submit_local_release(state, owner, vshard, reason).await;
};
let status_fn = match state.raft_status_fn.get() {
Some(f) => f,
None => {
tracing::warn!(
"release-reservation: raft status fn not installed; leaving release to lease GC"
);
return Ok(());
}
};
let leader = status_fn()
.into_iter()
.find(|g| g.group_id == SEQUENCER_GROUP_ID)
.map(|g| g.leader_id)
.unwrap_or(0);
if leader == 0 {
tracing::warn!(
"release-reservation: no sequencer leader elected yet; leaving release to lease GC"
);
return Ok(());
}
if leader == state.node_id {
return submit_local_release(state, owner, vshard, reason).await;
}
let mut targets = BTreeSet::new();
targets.insert(leader);
register_peers_from_topology(state, transport, &targets);
let owner_bytes = match zerompk::to_msgpack_vec(&owner) {
Ok(b) => b,
Err(e) => {
tracing::warn!(
"release-reservation: failed to encode TxnIdWire owner: {e}; leaving release to \
lease GC"
);
return Ok(());
}
};
let reason_bytes = match zerompk::to_msgpack_vec(&reason) {
Ok(b) => b,
Err(e) => {
tracing::warn!(
"release-reservation: failed to encode ReleaseReason: {e}; leaving release to \
lease GC"
);
return Ok(());
}
};
let deadline_remaining_ms = state
.tuning
.network
.default_deadline_secs
.saturating_mul(1000)
.max(1);
let req = ReleaseReservationRequest {
owner_bytes,
vshard,
reason_bytes,
deadline_remaining_ms,
trace_id: [0u8; 16],
};
let read_timeout = Duration::from_millis(deadline_remaining_ms.saturating_add(2_000));
match transport
.send_rpc_with_read_timeout(
leader,
RaftRpc::ReleaseReservationRequest(req),
read_timeout,
)
.await
{
Ok(RaftRpc::ReleaseReservationResponse(ReleaseReservationResponse { error: None })) => {
Ok(())
}
Ok(RaftRpc::ReleaseReservationResponse(ReleaseReservationResponse { error: Some(e) })) => {
tracing::warn!(
"release-reservation failed on sequencer leader node {leader}: {e:?}; leaving \
release to lease GC"
);
Ok(())
}
Ok(other) => {
tracing::warn!(
"release-reservation: unexpected reply from node {leader}: {other:?}; leaving \
release to lease GC"
);
Ok(())
}
Err(e) => {
tracing::warn!(
"release-reservation RPC to sequencer leader node {leader} failed: {e}; leaving \
release to lease GC"
);
Ok(())
}
}
}