pub mod abi;
pub mod capability;
pub mod codec;
pub mod contract;
pub mod error;
pub mod handle;
pub mod liveliness;
pub mod metadata;
pub mod query;
mod runtime_metrics;
pub mod server;
pub mod session;
pub mod topic;
pub use abi::{CodecId, encoding_string, parse_encoding_string};
pub use capability::OwnerCap;
pub use codec::{Codec, CodecError, MessagePack};
pub use contract::{ApiVersion, ContractBody, TopicRole};
pub use error::{BusError, Result};
pub use handle::{DEFAULT_QUERY_TIMEOUT, Latest, Publisher, Querier, Received, Subscriber};
pub use liveliness::{
ParticipantLivelinessEvent, ParticipantLivelinessKey, ParticipantLivelinessObserver,
ParticipantLivelinessStatus, ParticipantLivelinessToken,
};
pub use metadata::{BusMetadata, Source};
pub use query::{QueryCode, QueryError, QueryFailure, ServerResult};
#[doc(hidden)]
pub use runtime_metrics::{
RuntimeBufferKind, RuntimeDirection, RuntimeMetricKey, RuntimeMetricSnapshot,
};
pub use server::{IncomingQuery, ServerQueryable};
pub use session::{Bus, BusConfig, BusHealth};
pub use topic::{AskQuery, Publish, ServeQuery, Subscribe, Topic, TopicKind, WildcardPublish};
use std::collections::VecDeque;
#[doc(hidden)]
#[derive(Debug, Default)]
pub struct RetiredEpochs {
epochs: VecDeque<u64>,
}
impl RetiredEpochs {
pub const CAPACITY: usize = 8;
pub fn contains(&self, epoch: u64) -> bool {
self.epochs.contains(&epoch)
}
pub fn retire(&mut self, epoch: u64) {
if epoch == 0 {
return;
}
self.epochs.retain(|candidate| *candidate != epoch);
if self.epochs.len() == Self::CAPACITY {
self.epochs.pop_front();
}
self.epochs.push_back(epoch);
}
pub fn activate(&mut self, epoch: u64) {
self.epochs.retain(|candidate| *candidate != epoch);
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct LogicalTime {
epoch: u64,
time_ns: u64,
}
impl LogicalTime {
pub const fn new(epoch: u64, time_ns: u64) -> Self {
LogicalTime { epoch, time_ns }
}
pub const fn epoch(self) -> u64 {
self.epoch
}
pub const fn time_ns(self) -> u64 {
self.time_ns
}
}
#[cfg(test)]
mod tests;