#[cfg(rings_native)]
pub(crate) mod engine;
pub(crate) mod platform;
#[cfg(rings_browser)]
pub(crate) mod wt;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::Ordering;
#[cfg(rings_native)]
use std::time::Duration;
use bytes::Bytes;
use rings_core::dht::Did;
use serde::Deserialize;
use serde::Serialize;
#[cfg(rings_native)]
pub(crate) const RELAY_IDLE_TIMEOUT: Duration = Duration::from_secs(5 * 60);
pub(crate) fn allocate_non_reusing(counter: &AtomicU64) -> Option<u64> {
counter
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |value| {
value.checked_add(1)
})
.ok()
}
#[derive(Clone, Copy, PartialEq, Eq, Hash, Debug, Serialize, Deserialize)]
pub struct SessionId(pub u64);
#[derive(Clone, Copy, PartialEq, Eq, Hash, Debug, Serialize, Deserialize)]
pub enum Initiator {
Local,
Remote,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum EffectEnqueue {
Enqueued,
Missing,
Failed,
}
#[cfg(any(test, rings_browser))]
pub(crate) const MAX_OUTBOUND_QUEUE_OPS: usize = 1024;
#[cfg(any(test, rings_browser))]
pub(crate) const MAX_OUTBOUND_QUEUE_BYTES: usize = 8 * 1024 * 1024;
#[derive(Default)]
#[cfg(any(test, rings_browser))]
pub(crate) struct OutboundQueueBudget {
operations: usize,
data_bytes: usize,
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
#[cfg(any(test, rings_browser))]
pub(crate) enum OutboundDrainState {
#[default]
Idle,
Active,
}
#[cfg(any(test, rings_browser))]
impl OutboundDrainState {
pub(crate) fn claim(&mut self) -> bool {
match self {
Self::Idle => {
*self = Self::Active;
true
}
Self::Active => false,
}
}
pub(crate) fn release(&mut self) {
*self = Self::Idle;
}
}
#[cfg(any(test, rings_browser))]
impl OutboundQueueBudget {
pub(crate) fn try_reserve(&mut self, data_bytes: usize) -> bool {
let Some(operations) = self.operations.checked_add(1) else {
return false;
};
let Some(total_bytes) = self.data_bytes.checked_add(data_bytes) else {
return false;
};
if operations > MAX_OUTBOUND_QUEUE_OPS || total_bytes > MAX_OUTBOUND_QUEUE_BYTES {
return false;
}
self.operations = operations;
self.data_bytes = total_bytes;
true
}
pub(crate) fn release(&mut self, data_bytes: usize) -> bool {
let Some(operations) = self.operations.checked_sub(1) else {
return false;
};
let Some(total_bytes) = self.data_bytes.checked_sub(data_bytes) else {
return false;
};
self.operations = operations;
self.data_bytes = total_bytes;
true
}
}
#[derive(Clone, PartialEq, Eq, Hash, Debug)]
pub struct SessionKey {
pub peer: Did,
pub namespace: String,
pub session: SessionId,
pub initiator: Initiator,
}
impl SessionKey {
pub fn new(
peer: Did,
namespace: impl Into<String>,
session: SessionId,
initiator: Initiator,
) -> Self {
Self {
peer,
namespace: namespace.into(),
session,
initiator,
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum TransportKind {
Tcp,
Udp,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub enum Frame {
Open {
session: SessionId,
service: String,
},
Data {
session: SessionId,
from_opener: bool,
bytes: Bytes,
},
Shutdown {
session: SessionId,
from_opener: bool,
},
Close {
session: SessionId,
from_opener: bool,
},
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_outbound_budget_preserves_both_operation_and_byte_bounds() {
let mut operation_bound = OutboundQueueBudget::default();
for _ in 0..MAX_OUTBOUND_QUEUE_OPS {
assert!(operation_bound.try_reserve(0));
}
assert!(!operation_bound.try_reserve(0));
let mut byte_bound = OutboundQueueBudget::default();
assert!(byte_bound.try_reserve(MAX_OUTBOUND_QUEUE_BYTES));
assert!(!byte_bound.try_reserve(1));
}
#[test]
fn test_rejected_outbound_budget_reservation_does_not_consume_capacity() {
let mut budget = OutboundQueueBudget::default();
assert!(!budget.try_reserve(MAX_OUTBOUND_QUEUE_BYTES + 1));
assert!(budget.try_reserve(MAX_OUTBOUND_QUEUE_BYTES));
}
#[test]
fn test_released_outbound_budget_can_be_reserved_again() {
let mut budget = OutboundQueueBudget::default();
assert!(budget.try_reserve(MAX_OUTBOUND_QUEUE_BYTES));
assert!(budget.release(MAX_OUTBOUND_QUEUE_BYTES));
assert!(budget.try_reserve(MAX_OUTBOUND_QUEUE_BYTES));
}
#[test]
fn test_invalid_outbound_budget_release_is_total_and_does_not_mutate() {
let mut budget = OutboundQueueBudget::default();
assert!(budget.try_reserve(4));
assert!(!budget.release(5));
assert!(budget.release(4));
assert!(!budget.release(0));
}
#[test]
fn test_outbound_drain_has_exactly_one_owner_until_empty() {
let mut drain = OutboundDrainState::Idle;
assert!(drain.claim());
assert!(!drain.claim());
drain.release();
assert!(drain.claim());
}
#[test]
fn test_non_reusing_allocator_is_total_at_exhaustion() {
let counter = AtomicU64::new(u64::MAX - 1);
assert_eq!(allocate_non_reusing(&counter), Some(u64::MAX - 1));
assert_eq!(allocate_non_reusing(&counter), None);
assert_eq!(counter.load(Ordering::Relaxed), u64::MAX);
}
}