use std::sync::{
Arc,
atomic::{AtomicU64, Ordering},
};
struct AuthorityState {
next_request: AtomicU64,
}
#[doc(hidden)]
pub struct DbPhysicalDispositionIssuer {
state: Arc<AuthorityState>,
}
#[doc(hidden)]
pub struct DbPhysicalStartupHalf {
state: Arc<AuthorityState>,
}
#[doc(hidden)]
pub struct DbPhysicalRequestIssuer {
state: Arc<AuthorityState>,
}
#[doc(hidden)]
pub struct DbPhysicalProcessCapability {
state: Arc<AuthorityState>,
}
#[doc(hidden)]
pub struct DbPhysicalRequestHalf {
state: Arc<AuthorityState>,
request: u64,
}
#[doc(hidden)]
pub struct DbPhysicalExecutionHalf {
state: Arc<AuthorityState>,
request: u64,
}
#[doc(hidden)]
pub struct DbPhysicalDispositionOwner<T> {
state: Arc<AuthorityState>,
request: u64,
disposition: DbPhysicalDisposition,
value: T,
}
#[doc(hidden)]
pub struct DbPhysicalDispositionReceipt<T> {
disposition: DbPhysicalDisposition,
value: T,
}
#[doc(hidden)]
pub struct DbRequestNotUsedReceipt<T> {
value: T,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[doc(hidden)]
pub enum DbPhysicalDisposition {
Returned,
Discarded,
}
impl DbPhysicalDispositionIssuer {
#[doc(hidden)]
pub fn issue() -> Self {
Self {
state: Arc::new(AuthorityState {
next_request: AtomicU64::new(1),
}),
}
}
#[doc(hidden)]
pub fn into_startup_and_request_issuer(
self,
) -> (DbPhysicalStartupHalf, DbPhysicalRequestIssuer) {
(
DbPhysicalStartupHalf {
state: Arc::clone(&self.state),
},
DbPhysicalRequestIssuer { state: self.state },
)
}
}
impl DbPhysicalStartupHalf {
#[doc(hidden)]
pub fn into_process_capability(self) -> DbPhysicalProcessCapability {
DbPhysicalProcessCapability { state: self.state }
}
}
impl DbPhysicalRequestIssuer {
#[doc(hidden)]
pub fn issue_request(&self) -> Option<(DbPhysicalRequestHalf, DbPhysicalExecutionHalf)> {
let request = self
.state
.next_request
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
current.checked_add(1)
})
.ok()?;
Some((
DbPhysicalRequestHalf {
state: Arc::clone(&self.state),
request,
},
DbPhysicalExecutionHalf {
state: Arc::clone(&self.state),
request,
},
))
}
}
impl DbPhysicalProcessCapability {
#[doc(hidden)]
pub fn connection_returned<T>(
&self,
execution: DbPhysicalExecutionHalf,
value: T,
) -> Result<DbPhysicalDispositionOwner<T>, (DbPhysicalExecutionHalf, T)> {
self.seal(execution, DbPhysicalDisposition::Returned, value)
}
#[doc(hidden)]
pub fn connection_discarded<T>(
&self,
execution: DbPhysicalExecutionHalf,
value: T,
) -> Result<DbPhysicalDispositionOwner<T>, (DbPhysicalExecutionHalf, T)> {
self.seal(execution, DbPhysicalDisposition::Discarded, value)
}
fn seal<T>(
&self,
execution: DbPhysicalExecutionHalf,
disposition: DbPhysicalDisposition,
value: T,
) -> Result<DbPhysicalDispositionOwner<T>, (DbPhysicalExecutionHalf, T)> {
if !Arc::ptr_eq(&self.state, &execution.state) {
return Err((execution, value));
}
Ok(DbPhysicalDispositionOwner {
state: execution.state,
request: execution.request,
disposition,
value,
})
}
}
#[doc(hidden)]
pub fn pair_db_physical_disposition<T>(
physical: DbPhysicalDispositionOwner<T>,
request: DbPhysicalRequestHalf,
) -> Result<DbPhysicalDispositionReceipt<T>, (DbPhysicalDispositionOwner<T>, DbPhysicalRequestHalf)>
{
if !Arc::ptr_eq(&physical.state, &request.state) || physical.request != request.request {
return Err((physical, request));
}
Ok(DbPhysicalDispositionReceipt {
disposition: physical.disposition,
value: physical.value,
})
}
impl<T> DbPhysicalDispositionReceipt<T> {
#[doc(hidden)]
pub fn into_outcome(self) -> (DbPhysicalDisposition, T) {
(self.disposition, self.value)
}
}
#[doc(hidden)]
pub fn seal_db_request_not_used<T>(
request: DbPhysicalRequestHalf,
execution: DbPhysicalExecutionHalf,
value: T,
) -> Result<DbRequestNotUsedReceipt<T>, (DbPhysicalRequestHalf, DbPhysicalExecutionHalf, T)> {
if !Arc::ptr_eq(&request.state, &execution.state) || request.request != execution.request {
return Err((request, execution, value));
}
Ok(DbRequestNotUsedReceipt { value })
}
impl<T> DbRequestNotUsedReceipt<T> {
#[doc(hidden)]
pub fn into_value(self) -> T {
self.value
}
}
#[cfg(test)]
mod tests {
use super::*;
fn authority() -> (DbPhysicalProcessCapability, DbPhysicalRequestIssuer) {
let (startup, requests) =
DbPhysicalDispositionIssuer::issue().into_startup_and_request_issuer();
(startup.into_process_capability(), requests)
}
#[test]
fn same_process_and_request_seal_returned_and_discarded() {
let (process, requests) = authority();
let (request, execution) = requests.issue_request().unwrap();
let physical = match process.connection_returned(execution, "query") {
Ok(value) => value,
Err(_) => panic!("same process rejected"),
};
let receipt = match pair_db_physical_disposition(physical, request) {
Ok(value) => value,
Err(_) => panic!("same request rejected"),
};
assert_eq!(
receipt.into_outcome(),
(DbPhysicalDisposition::Returned, "query")
);
let (request, execution) = requests.issue_request().unwrap();
let physical = match process.connection_discarded(execution, "write") {
Ok(value) => value,
Err(_) => panic!("same process rejected"),
};
let receipt = match pair_db_physical_disposition(physical, request) {
Ok(value) => value,
Err(_) => panic!("same request rejected"),
};
assert_eq!(
receipt.into_outcome(),
(DbPhysicalDisposition::Discarded, "write")
);
}
#[test]
fn foreign_process_and_crossed_request_return_all_owners() {
let (first_process, first_requests) = authority();
let (second_process, second_requests) = authority();
let (first_request, first_execution) = first_requests.issue_request().unwrap();
let (second_request, second_execution) = second_requests.issue_request().unwrap();
let (first_execution, value) = match second_process.connection_returned(first_execution, 11)
{
Ok(_) => panic!("foreign process accepted"),
Err(owners) => owners,
};
let first_physical = match first_process.connection_returned(first_execution, value) {
Ok(value) => value,
Err(_) => panic!("original process rejected"),
};
let second_physical = match second_process.connection_discarded(second_execution, 22) {
Ok(value) => value,
Err(_) => panic!("original process rejected"),
};
let (first_physical, second_request) =
match pair_db_physical_disposition(first_physical, second_request) {
Ok(_) => panic!("crossed request accepted"),
Err(owners) => owners,
};
let (second_physical, first_request) =
match pair_db_physical_disposition(second_physical, first_request) {
Ok(_) => panic!("crossed request accepted"),
Err(owners) => owners,
};
let first = match pair_db_physical_disposition(first_physical, first_request) {
Ok(value) => value,
Err(_) => panic!("original request rejected"),
};
assert_eq!(first.into_outcome(), (DbPhysicalDisposition::Returned, 11));
let second = match pair_db_physical_disposition(second_physical, second_request) {
Ok(value) => value,
Err(_) => panic!("original request rejected"),
};
assert_eq!(
second.into_outcome(),
(DbPhysicalDisposition::Discarded, 22)
);
}
#[test]
fn untouched_request_seals_not_used_and_crossed_halves_retry() {
let (_first_process, first_requests) = authority();
let (_second_process, second_requests) = authority();
let (first_request, first_execution) = first_requests.issue_request().unwrap();
let (second_request, second_execution) = second_requests.issue_request().unwrap();
let (first_request, second_execution, first_value) =
match seal_db_request_not_used(first_request, second_execution, "first") {
Ok(_) => panic!("crossed request accepted"),
Err(owners) => owners,
};
let (second_request, first_execution, second_value) =
match seal_db_request_not_used(second_request, first_execution, "second") {
Ok(_) => panic!("crossed request accepted"),
Err(owners) => owners,
};
let first = seal_db_request_not_used(first_request, first_execution, first_value)
.unwrap_or_else(|_| panic!("original request rejected"));
let second = seal_db_request_not_used(second_request, second_execution, second_value)
.unwrap_or_else(|_| panic!("original request rejected"));
assert_eq!(first.into_value(), "first");
assert_eq!(second.into_value(), "second");
}
}