use std::{
io,
time::{Duration, Instant},
};
use thiserror::Error;
use crate::snapshot::SnapshotReadLimits;
#[derive(Clone, Debug, Default, Eq, PartialEq)]
pub struct StorageLimits {
pub recovery: RecoveryLimits,
pub maintenance: MaintenanceLimits,
}
impl StorageLimits {
pub(crate) fn compatibility() -> Self {
Self {
recovery: RecoveryLimits::compatibility(),
maintenance: MaintenanceLimits::compatibility(),
}
}
pub(crate) fn validate(&self) -> Result<(), StorageLimitError> {
self.recovery.validate()?;
self.maintenance.validate()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RecoveryLimits {
pub timeout: Duration,
pub max_directory_entries: u64,
pub max_log_file_bytes: u64,
pub max_log_frames: u64,
pub max_transactions: u64,
pub max_operations: u64,
pub max_decoded_operation_bytes: u64,
pub snapshot: SnapshotReadLimits,
pub max_lexical_documents: u64,
pub max_lexical_tokens: u64,
}
impl Default for RecoveryLimits {
fn default() -> Self {
Self {
timeout: Duration::from_secs(60),
max_directory_entries: 1_000_000,
max_log_file_bytes: 2 * 1024 * 1024 * 1024,
max_log_frames: 1_000_000,
max_transactions: 1_000_000,
max_operations: 1_000_000,
max_decoded_operation_bytes: 1024 * 1024 * 1024,
snapshot: SnapshotReadLimits::default(),
max_lexical_documents: 1_000_000,
max_lexical_tokens: 10_000_000,
}
}
}
impl RecoveryLimits {
pub(crate) fn compatibility() -> Self {
Self {
timeout: Duration::MAX,
max_directory_entries: u64::MAX,
max_log_file_bytes: u64::MAX,
max_log_frames: u64::MAX,
max_transactions: u64::MAX,
max_operations: u64::MAX,
max_decoded_operation_bytes: u64::MAX,
snapshot: SnapshotReadLimits {
file_bytes: u64::MAX,
entries: u64::MAX,
decoded_bytes: u64::MAX,
},
max_lexical_documents: u64::MAX,
max_lexical_tokens: u64::MAX,
}
}
pub(crate) fn validate(&self) -> Result<(), StorageLimitError> {
require_nonzero_duration(self.timeout, "recovery.timeout")?;
for (name, value) in [
("recovery.max_directory_entries", self.max_directory_entries),
("recovery.max_log_file_bytes", self.max_log_file_bytes),
("recovery.max_log_frames", self.max_log_frames),
("recovery.max_transactions", self.max_transactions),
("recovery.max_operations", self.max_operations),
(
"recovery.max_decoded_operation_bytes",
self.max_decoded_operation_bytes,
),
("recovery.snapshot.file_bytes", self.snapshot.file_bytes),
("recovery.snapshot.entries", self.snapshot.entries),
(
"recovery.snapshot.decoded_bytes",
self.snapshot.decoded_bytes,
),
("recovery.max_lexical_documents", self.max_lexical_documents),
("recovery.max_lexical_tokens", self.max_lexical_tokens),
] {
require_nonzero(value, name)?;
}
Ok(())
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct MaintenanceLimits {
pub timeout: Duration,
pub snapshot: SnapshotReadLimits,
}
impl Default for MaintenanceLimits {
fn default() -> Self {
Self {
timeout: Duration::from_secs(60),
snapshot: SnapshotReadLimits::default(),
}
}
}
impl MaintenanceLimits {
fn compatibility() -> Self {
Self {
timeout: Duration::MAX,
snapshot: SnapshotReadLimits {
file_bytes: u64::MAX,
entries: u64::MAX,
decoded_bytes: u64::MAX,
},
}
}
pub(crate) fn validate(&self) -> Result<(), StorageLimitError> {
require_nonzero_duration(self.timeout, "maintenance.timeout")?;
for (name, value) in [
("maintenance.snapshot.file_bytes", self.snapshot.file_bytes),
("maintenance.snapshot.entries", self.snapshot.entries),
(
"maintenance.snapshot.decoded_bytes",
self.snapshot.decoded_bytes,
),
] {
require_nonzero(value, name)?;
}
Ok(())
}
}
#[derive(Clone, Debug, Error, Eq, PartialEq)]
pub enum StorageLimitError {
#[error("storage limit must be positive: {name}")]
ZeroLimit {
name: &'static str,
},
#[error("storage operation timed out")]
TimedOut,
#[error("storage directory entry limit exceeded: {maximum}")]
DirectoryEntriesExceeded {
maximum: u64,
},
#[error("log file byte limit exceeded: {actual} > {maximum}")]
LogFileBytesExceeded {
actual: u64,
maximum: u64,
},
#[error("log frame limit exceeded: {maximum}")]
LogFramesExceeded {
maximum: u64,
},
#[error("recovery transaction limit exceeded: {maximum}")]
TransactionsExceeded {
maximum: u64,
},
#[error("recovery operation limit exceeded: {maximum}")]
OperationsExceeded {
maximum: u64,
},
#[error("recovery decoded operation byte limit exceeded: {maximum}")]
DecodedOperationBytesExceeded {
maximum: u64,
},
#[error("lexical rebuild document limit exceeded: {maximum}")]
LexicalDocumentsExceeded {
maximum: u64,
},
#[error("lexical rebuild token limit exceeded: {maximum}")]
LexicalTokensExceeded {
maximum: u64,
},
}
#[derive(Clone, Debug)]
pub(crate) struct OperationDeadline {
started: Instant,
timeout: Duration,
}
impl OperationDeadline {
pub(crate) fn new(timeout: Duration) -> Self {
Self {
started: Instant::now(),
timeout,
}
}
pub(crate) fn check(&self) -> Result<(), StorageLimitError> {
if self.started.elapsed() >= self.timeout {
Err(StorageLimitError::TimedOut)
} else {
Ok(())
}
}
}
pub(crate) fn limit_io_error(source: StorageLimitError) -> io::Error {
let kind = if matches!(source, StorageLimitError::TimedOut) {
io::ErrorKind::TimedOut
} else {
io::ErrorKind::Other
};
io::Error::new(kind, source)
}
pub fn storage_limit_from_io(source: &io::Error) -> Option<&StorageLimitError> {
source
.get_ref()
.and_then(|source| source.downcast_ref::<StorageLimitError>())
}
fn require_nonzero(value: u64, name: &'static str) -> Result<(), StorageLimitError> {
if value == 0 {
Err(StorageLimitError::ZeroLimit { name })
} else {
Ok(())
}
}
fn require_nonzero_duration(value: Duration, name: &'static str) -> Result<(), StorageLimitError> {
if value.is_zero() {
Err(StorageLimitError::ZeroLimit { name })
} else {
Ok(())
}
}