use std::io;
use std::sync::Arc;
use std::sync::atomic::{AtomicU8, Ordering};
use std::thread::{self, Thread};
use crate::sync::Mutex;
use super::request::WriteRequest;
const SLOT_IDLE: u8 = 0;
const SLOT_PENDING: u8 = 1;
const SLOT_DONE: u8 = 2;
struct SlotPayload {
request: WriteRequest,
outcome: Option<io::Result<u64>>,
}
pub(crate) struct WriteSlot {
state: AtomicU8,
payload: Mutex<SlotPayload>,
thread: Thread,
}
impl WriteSlot {
pub(super) fn new() -> Self {
Self {
state: AtomicU8::new(SLOT_IDLE),
payload: Mutex::new(SlotPayload {
request: WriteRequest::Idle,
outcome: None,
}),
thread: thread::current(),
}
}
pub(super) fn arm(&self, request: WriteRequest) -> Result<(), WriteRequest> {
if self.state.load(Ordering::Acquire) != SLOT_IDLE {
return Err(request);
}
{
let mut payload = self.payload.lock();
payload.request = request;
payload.outcome = None;
}
self.state.store(SLOT_PENDING, Ordering::Release);
Ok(())
}
pub(super) fn take_request(&self) -> WriteRequest {
std::mem::replace(&mut self.payload.lock().request, WriteRequest::Idle)
}
pub(super) fn complete(&self, outcome: io::Result<u64>) {
self.payload.lock().outcome = Some(outcome);
self.state.store(SLOT_DONE, Ordering::Release);
self.thread.unpark();
}
pub(super) fn is_done(&self) -> bool {
self.state.load(Ordering::Acquire) == SLOT_DONE
}
pub(super) fn finish(&self) -> io::Result<u64> {
let outcome = self.payload.lock().outcome.take();
self.state.store(SLOT_IDLE, Ordering::Release);
outcome.unwrap_or_else(|| Err(io::Error::other("commit slot completed with no outcome")))
}
}
thread_local! {
static WRITE_SLOT: Arc<WriteSlot> = Arc::new(WriteSlot::new());
}
pub(super) fn thread_slot() -> Arc<WriteSlot> {
WRITE_SLOT.with(Arc::clone)
}
#[cfg(test)]
mod tests {
use super::super::super::DurabilityMode;
use super::*;
fn put(key: &[u8], value: &[u8]) -> WriteRequest {
WriteRequest::Put {
key: key.to_vec(),
value: value.to_vec(),
durability: DurabilityMode::Eventual,
disable_wal: false,
}
}
#[test]
fn slot_round_trips_a_request_and_an_outcome() {
let slot = WriteSlot::new();
assert!(slot.arm(put(b"k", b"v")).is_ok());
assert!(!slot.is_done());
let taken = slot.take_request();
assert_eq!(taken.op_count(), 1);
slot.complete(Ok(7));
assert!(slot.is_done());
assert!(slot.finish().is_ok());
assert!(slot.arm(put(b"k2", b"v2")).is_ok());
}
#[test]
fn arm_refuses_a_slot_with_an_outstanding_ticket() {
let slot = WriteSlot::new();
assert!(slot.arm(put(b"k", b"v")).is_ok());
match slot.arm(put(b"k2", b"v2")) {
Err(returned) => assert_eq!(returned.op_count(), 1),
Ok(()) => panic!("a pending slot must refuse a second request"),
}
}
#[test]
fn a_completed_slot_with_no_outcome_fails_loud() {
let slot = WriteSlot::new();
slot.state.store(SLOT_DONE, Ordering::Release);
let err = slot.finish().expect_err("missing outcome must be an error");
assert!(err.to_string().contains("no outcome"));
}
}