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,
scope: u64,
observation_taken: bool,
}
#[doc(hidden)]
pub struct DbScopeObservation<C> {
state: Arc<AuthorityState>,
request: u64,
scope: u64,
context: C,
}
#[doc(hidden)]
pub struct DbScopeTerminalObservation<C> {
context: C,
fields: DbScopeLogFields,
}
#[derive(serde::Serialize)]
#[doc(hidden)]
pub struct DbScopeLogFields {
transaction_scope: u64,
}
impl DbPhysicalRequestHalf {
#[doc(hidden)]
pub fn take_scope_observation<C>(
&mut self,
execution: &DbPhysicalExecutionHalf,
context: C,
) -> Result<DbScopeObservation<C>, C> {
if self.observation_taken
|| !Arc::ptr_eq(&self.state, &execution.state)
|| self.request != execution.request
|| self.scope != execution.scope
{
return Err(context);
}
self.observation_taken = true;
Ok(DbScopeObservation {
state: Arc::clone(&self.state),
request: self.request,
scope: self.scope,
context,
})
}
}
impl<C> DbScopeObservation<C> {
#[doc(hidden)]
pub fn bind_terminal(
self,
execution: &DbPhysicalExecutionHalf,
) -> Result<DbScopeTerminalObservation<C>, Self> {
if !Arc::ptr_eq(&self.state, &execution.state)
|| self.request != execution.request
|| self.scope != execution.scope
{
return Err(self);
}
Ok(DbScopeTerminalObservation {
context: self.context,
fields: DbScopeLogFields {
transaction_scope: self.scope,
},
})
}
}
impl<C> DbScopeTerminalObservation<C> {
#[doc(hidden)]
pub fn into_log_parts(self) -> (C, DbScopeLogFields) {
(self.context, self.fields)
}
}
#[doc(hidden)]
pub struct DbPhysicalExecutionHalf {
state: Arc<AuthorityState>,
request: u64,
scope: u64,
}
#[doc(hidden)]
pub struct DbPhysicalDispositionOwner<T> {
state: Arc<AuthorityState>,
request: u64,
scope: u64,
disposition: DbPhysicalDisposition,
value: T,
}
#[doc(hidden)]
pub struct DbPhysicalDispositionReceipt<T> {
request: DbPhysicalRequestHalf,
disposition: DbPhysicalDisposition,
value: T,
}
#[doc(hidden)]
pub struct DbScopeContinuation {
request: DbPhysicalRequestHalf,
}
impl DbScopeContinuation {
pub fn into_next_scope(self) -> Result<(DbPhysicalRequestHalf, DbPhysicalExecutionHalf), Self> {
let Some(scope) = self.request.scope.checked_add(1) else {
return Err(self);
};
let execution = DbPhysicalExecutionHalf {
state: Arc::clone(&self.request.state),
request: self.request.request,
scope,
};
Ok((
DbPhysicalRequestHalf {
scope,
observation_taken: false,
..self.request
},
execution,
))
}
}
#[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,
scope: 0,
observation_taken: false,
},
DbPhysicalExecutionHalf {
state: Arc::clone(&self.state),
request,
scope: 0,
},
))
}
}
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,
scope: execution.scope,
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
|| physical.scope != request.scope
{
return Err((physical, request));
}
Ok(DbPhysicalDispositionReceipt {
request,
disposition: physical.disposition,
value: physical.value,
})
}
impl<T> DbPhysicalDispositionReceipt<T> {
pub fn into_scope_continuation(self) -> (DbScopeContinuation, DbPhysicalDisposition, T) {
(
DbScopeContinuation {
request: self.request,
},
self.disposition,
self.value,
)
}
#[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
|| request.scope != execution.scope
{
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::*;
#[test]
fn scope_observation_foreign_retry_once_and_serial_scope() {
let (process, requests) = authority();
let (mut a, mut ea) = requests.issue_request().unwrap();
let (_, eb) = requests.issue_request().unwrap();
let context = a
.take_scope_observation(&eb, String::from("safe-context"))
.err()
.unwrap();
assert!(!a.observation_taken);
let token = a.take_scope_observation(&ea, context).ok().unwrap();
assert!(a.take_scope_observation(&ea, "repeat").is_err());
let token = token.bind_terminal(&eb).err().unwrap();
ea.scope = 1;
let token = token.bind_terminal(&ea).err().unwrap();
ea.scope = 0;
let (context, fields) = token.bind_terminal(&ea).ok().unwrap().into_log_parts();
assert_eq!(context, "safe-context");
assert_eq!(
serde_json::to_value(fields).unwrap(),
serde_json::json!({"transaction_scope":0})
);
let proof = process.connection_returned(ea, ()).ok().unwrap();
let (continuation, _, _) = pair_db_physical_disposition(proof, a)
.ok()
.unwrap()
.into_scope_continuation();
let (mut a, ea) = continuation.into_next_scope().ok().unwrap();
let token = a.take_scope_observation(&ea, context).ok().unwrap();
let (_, fields) = token.bind_terminal(&ea).ok().unwrap().into_log_parts();
assert_eq!(
serde_json::to_value(fields).unwrap(),
serde_json::json!({"transaction_scope":1})
);
assert_eq!(requests.state.next_request.load(Ordering::Relaxed), 3);
}
#[test]
fn serial_scope_identity_cross_and_drift_restore() {
let (process, issuer) = authority();
let (mut request, execution) = issuer.issue_request().unwrap();
let request_id = request.request;
let physical = process.connection_returned(execution, 7).ok().unwrap();
request.scope = 1; let (physical, mut request) = pair_db_physical_disposition(physical, request)
.err()
.unwrap();
request.scope = 0;
let receipt = pair_db_physical_disposition(physical, request)
.ok()
.unwrap();
let (continuation, disposition, value) = receipt.into_scope_continuation();
assert_eq!(disposition, DbPhysicalDisposition::Returned);
assert_eq!(value, 7);
let (request, execution) = continuation.into_next_scope().ok().unwrap();
assert_eq!(request.request, request_id);
assert_eq!(request.scope, 1);
assert_eq!(issuer.state.next_request.load(Ordering::Relaxed), 2);
let (foreign, _) = issuer.issue_request().unwrap();
let physical = process.connection_discarded(execution, 9).ok().unwrap();
let (physical, foreign) = pair_db_physical_disposition(physical, foreign)
.err()
.unwrap();
assert_ne!(foreign.request, request_id);
let receipt = pair_db_physical_disposition(physical, request)
.ok()
.unwrap();
assert_eq!(
receipt.into_outcome(),
(DbPhysicalDisposition::Discarded, 9)
);
}
#[test]
fn serial_scope_overflow_returns_same_continuation() {
let (_, issuer) = authority();
let (mut request, _) = issuer.issue_request().unwrap();
request.scope = u64::MAX;
let continuation = DbScopeContinuation { request };
let continuation = continuation.into_next_scope().err().unwrap();
assert_eq!(continuation.request.scope, u64::MAX);
assert_eq!(continuation.request.request, 1);
}
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");
}
}