use tokio::sync::mpsc;
use tokio::sync::oneshot;
use tracing::warn;
use crate::calvin::sequencer::error::SequencerError;
use crate::calvin::types::{LockKeyWire, ReleaseReason, TxnIdWire};
pub enum ReservationRequest {
Reserve {
key: LockKeyWire,
vshard: u32,
owner: Option<TxnIdWire>,
reply: oneshot::Sender<TxnIdWire>,
},
Release {
owner: TxnIdWire,
vshard: u32,
reason: ReleaseReason,
},
}
#[derive(Clone)]
pub struct ReservationInbox {
tx: mpsc::Sender<ReservationRequest>,
}
pub struct ReservationInboxReceiver {
rx: mpsc::Receiver<ReservationRequest>,
}
impl ReservationInbox {
pub fn submit_reserve(
&self,
key: LockKeyWire,
vshard: u32,
owner: Option<TxnIdWire>,
) -> Result<oneshot::Receiver<TxnIdWire>, SequencerError> {
let (reply, reply_rx) = oneshot::channel();
let request = ReservationRequest::Reserve {
key,
vshard,
owner,
reply,
};
match self.tx.try_send(request) {
Ok(()) => Ok(reply_rx),
Err(mpsc::error::TrySendError::Full(_)) => Err(SequencerError::Overloaded),
Err(mpsc::error::TrySendError::Closed(_)) => {
warn!("reservation inbox channel closed; sequencer service may have exited");
Err(SequencerError::Overloaded)
}
}
}
pub fn submit_release(
&self,
owner: TxnIdWire,
vshard: u32,
reason: ReleaseReason,
) -> Result<(), SequencerError> {
let request = ReservationRequest::Release {
owner,
vshard,
reason,
};
match self.tx.try_send(request) {
Ok(()) => Ok(()),
Err(mpsc::error::TrySendError::Full(_)) => Err(SequencerError::Overloaded),
Err(mpsc::error::TrySendError::Closed(_)) => {
warn!("reservation inbox channel closed; sequencer service may have exited");
Err(SequencerError::Overloaded)
}
}
}
}
impl ReservationInboxReceiver {
pub fn drain_into(&mut self, out: &mut Vec<ReservationRequest>) -> usize {
let mut count = 0;
while let Ok(request) = self.rx.try_recv() {
out.push(request);
count += 1;
}
count
}
pub fn drain_all_discard(&mut self) -> usize {
let mut count = 0;
while self.rx.try_recv().is_ok() {
count += 1;
}
count
}
}
pub fn new_reservation_inbox(capacity: usize) -> (ReservationInbox, ReservationInboxReceiver) {
let (tx, rx) = mpsc::channel(capacity);
(ReservationInbox { tx }, ReservationInboxReceiver { rx })
}
#[cfg(test)]
mod tests {
use super::*;
fn sample_key() -> LockKeyWire {
LockKeyWire::Kv {
collection: "sessions".to_owned(),
key: b"hot".to_vec(),
}
}
#[test]
fn submit_reserve_enqueues_and_returns_receiver() {
let (inbox, mut rx) = new_reservation_inbox(4);
let _reply_rx = inbox
.submit_reserve(sample_key(), 7, None)
.expect("submit should succeed");
let mut out = Vec::new();
let n = rx.drain_into(&mut out);
assert_eq!(n, 1);
assert!(matches!(
out[0],
ReservationRequest::Reserve {
vshard: 7,
owner: None,
..
}
));
}
#[test]
fn submit_reserve_returns_overloaded_when_full() {
let (inbox, _rx) = new_reservation_inbox(1);
let _rx = inbox
.submit_reserve(sample_key(), 1, None)
.expect("first submit fills channel");
let err = inbox
.submit_reserve(sample_key(), 1, None)
.expect_err("second submit should be overloaded");
assert_eq!(err, SequencerError::Overloaded);
}
#[test]
fn submit_release_is_fire_and_forget() {
let (inbox, mut rx) = new_reservation_inbox(4);
inbox
.submit_release(
TxnIdWire {
epoch: 3,
position: 9,
},
2,
ReleaseReason::Commit,
)
.expect("release should enqueue");
let mut out = Vec::new();
assert_eq!(rx.drain_into(&mut out), 1);
assert!(matches!(
out[0],
ReservationRequest::Release {
vshard: 2,
reason: ReleaseReason::Commit,
..
}
));
}
#[test]
fn drain_all_discard_drops_reply_senders() {
let (inbox, mut rx) = new_reservation_inbox(4);
let reply_rx = inbox.submit_reserve(sample_key(), 1, None).expect("submit");
let discarded = rx.drain_all_discard();
assert_eq!(discarded, 1);
assert!(reply_rx.blocking_recv().is_err());
}
#[test]
fn submit_reserve_echoes_existing_owner() {
let (inbox, mut rx) = new_reservation_inbox(4);
let existing = TxnIdWire {
epoch: 5,
position: 1 << 31,
};
inbox
.submit_reserve(sample_key(), 4, Some(existing))
.expect("submit");
let mut out = Vec::new();
rx.drain_into(&mut out);
match &out[0] {
ReservationRequest::Reserve { owner, .. } => {
assert_eq!(*owner, Some(existing));
}
_ => panic!("expected Reserve"),
}
}
}