pub mod async_api;
#[cfg(feature = "tokio")]
pub mod async_tokio;
pub mod backend;
pub mod cgroup;
pub mod error;
pub mod flow;
pub mod lease;
pub mod observability;
pub mod operation;
pub mod raw_api;
pub mod result;
pub mod rings;
pub mod torus_panic;
pub use error::{Error, Result};
pub use flow::Flow;
pub use lease::{LeaseError, LeaseRegistry, SharedLeaseRegistry};
pub use operation::{IoSlice, Operation};
pub use result::Result as TorusResult;
pub use rings::{CompletionRing, SubmissionRing};
pub use torus_panic::TorusPanic;
use backend::Backend;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::{Arc, Mutex};
pub struct Torus {
sq: SubmissionRing,
cq: CompletionRing,
backend: Mutex<Box<dyn Backend>>,
#[cfg(feature = "tracing")]
spans: Mutex<std::collections::HashMap<u64, crate::observability::FlowSpan>>,
}
unsafe impl Send for Torus {}
unsafe impl Sync for Torus {}
impl Torus {
pub fn new(ring_entries: u32, backend: Box<dyn Backend>) -> Result<Self> {
if !ring_entries.is_power_of_two() {
return Err(Error::InvalidParam("ring_entries must be a power of two"));
}
Ok(Self {
sq: SubmissionRing::new(ring_entries),
cq: CompletionRing::new(ring_entries),
backend: Mutex::new(backend),
#[cfg(feature = "tracing")]
spans: Mutex::new(std::collections::HashMap::new()),
})
}
pub fn submit(&self, flow: &Flow) -> Result<()> {
let n = self
.backend
.lock()
.unwrap()
.submit(std::slice::from_ref(flow))?;
if n == 0 {
return Err(Error::SubmissionFull);
}
#[cfg(feature = "tracing")]
{
self.spans
.lock()
.unwrap()
.insert(flow.user_data(), crate::observability::FlowSpan::new(flow));
}
Ok(())
}
pub fn submit_batch(&self, flows: &[Flow]) -> Result<usize> {
let n = self.backend.lock().unwrap().submit(flows)?;
#[cfg(feature = "tracing")]
{
let mut spans = self.spans.lock().unwrap();
for flow in flows.iter().take(n) {
spans.insert(flow.user_data(), crate::observability::FlowSpan::new(flow));
}
}
Ok(n)
}
pub fn readv(&self, fd: i32, bufs: &[IoSlice], offset: u64, user_data: u64) -> Result<()> {
let flow = Flow::with_user_data(
Operation::Readv {
fd,
bufs: bufs.as_ptr(),
buf_count: bufs.len() as u32,
offset,
},
user_data,
);
self.submit(&flow)
}
pub fn writev(&self, fd: i32, bufs: &[IoSlice], offset: u64, user_data: u64) -> Result<()> {
let flow = Flow::with_user_data(
Operation::Writev {
fd,
bufs: bufs.as_ptr(),
buf_count: bufs.len() as u32,
offset,
},
user_data,
);
self.submit(&flow)
}
pub fn reap(&self, results: &mut Vec<TorusResult>) -> Result<usize> {
let count = self.backend.lock().unwrap().reap(results)?;
#[cfg(feature = "tracing")]
{
let mut spans = self.spans.lock().unwrap();
for r in results.iter() {
if let Some(span) = spans.remove(&r.user_data) {
span.complete(r.result);
}
}
}
Ok(count)
}
pub fn wait(&self, timeout_us: u64) -> Result<()> {
#[cfg(feature = "tracing")]
let _span = tracing::debug_span!("torus_wait", timeout_us = timeout_us).entered();
self.backend.lock().unwrap().wait(timeout_us)
}
pub fn in_flight(&self) -> u32 {
self.backend.lock().unwrap().in_flight()
}
pub fn submission_ring(&self) -> &SubmissionRing {
&self.sq
}
pub fn completion_ring(&self) -> &CompletionRing {
&self.cq
}
#[cfg(unix)]
pub fn register_leases(&self, registry: &LeaseRegistry) -> crate::error::Result<()> {
let buffers = registry.as_register_buffers();
if buffers.is_empty() {
return Ok(());
}
self.backend.lock().unwrap().register_buffers(&buffers)
}
#[cfg(not(unix))]
pub fn register_leases(&self, _registry: &LeaseRegistry) -> crate::error::Result<()> {
Ok(())
}
pub unsafe fn raw(&self) -> raw_api::RawTorus<'_> {
raw_api::RawTorus::new(self)
}
}
pub type SharedTorus = Arc<Torus>;
pub struct TorusPool {
instances: Vec<Arc<Torus>>,
next: AtomicU32,
}
impl TorusPool {
pub fn new<F>(count: usize, ring_entries: u32, make_backend: F) -> Result<Self>
where
F: Fn(u32) -> Result<Box<dyn Backend>>,
{
if count == 0 {
return Err(Error::InvalidParam("pool count must be > 0"));
}
let mut instances = Vec::with_capacity(count);
for _ in 0..count {
let backend = make_backend(ring_entries)?;
instances.push(Arc::new(Torus::new(ring_entries, backend)?));
}
Ok(Self {
instances,
next: AtomicU32::new(0),
})
}
pub fn submit(&self, flow: &Flow) -> Result<()> {
let idx = self.next.fetch_add(1, Ordering::Relaxed) as usize % self.instances.len();
self.instances[idx].submit(flow)
}
pub fn submit_batch(&self, flows: &[Flow]) -> Result<usize> {
let mut total = 0;
for flow in flows {
self.submit(flow)?;
total += 1;
}
Ok(total)
}
pub fn reap(&self, results: &mut Vec<TorusResult>) -> Result<usize> {
let mut total = 0;
for instance in &self.instances {
total += instance.reap(results)?;
}
Ok(total)
}
pub fn len(&self) -> usize {
self.instances.len()
}
pub fn is_empty(&self) -> bool {
self.instances.is_empty()
}
pub fn get(&self, index: usize) -> Option<&Torus> {
self.instances.get(index).map(|arc| arc.as_ref())
}
pub fn in_flight(&self) -> u32 {
self.instances.iter().map(|i| i.in_flight()).sum()
}
}