mod auto_saved;
mod dataset;
mod handle;
mod kvs;
mod queue;
pub mod rate_limit;
pub use auto_saved::AutoSaved;
pub use dataset::{Dataset, DatasetExt, DatasetInfo, ListOptions, Page};
pub use handle::StorageHandle;
pub use kvs::{KeyInfo, KeyList, KeyValueStore, KeyValueStoreExt, KvEntry, ListKeysOptions};
pub use queue::{
AddOptions, AddRequestsBatchedResult, BatchAddHandle, Lease, LeaseId, ProcessedRequest,
QueueOpInfo, ReclaimOptions, RequestQueue, RequestSource,
};
pub use rate_limit::RateLimitReportingClient;
use std::sync::Arc;
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum StorageError {
#[error("serialization: {0}")]
Serialization(#[from] serde_json::Error),
#[error("io: {0}")]
Io(#[from] std::io::Error),
#[error("lease {lease_id} not found (already completed, abandoned, or expired)")]
LeaseNotFound {
lease_id: LeaseId,
},
#[error("storage backend: {0}")]
Backend(#[source] anyhow::Error),
#[error("operation not supported by this backend: {0}")]
Unsupported(&'static str),
#[error("storage backend rate limited")]
RateLimited {
retry_after: Option<std::time::Duration>,
},
}
impl StorageError {
pub fn is_rate_limited(&self) -> bool {
matches!(self, StorageError::RateLimited { .. })
}
}
pub type StorageResult<T> = Result<T, StorageError>;
impl From<StorageError> for crate::errors::CrawlError {
fn from(error: StorageError) -> Self {
match error {
StorageError::Serialization(error) => Self::NonRetryable(anyhow::Error::new(error)),
StorageError::Io(error) => Self::Retry(anyhow::Error::new(error)),
error @ StorageError::LeaseNotFound { .. } => {
Self::NonRetryable(anyhow::Error::new(error))
}
StorageError::Backend(error) => Self::Retry(error),
error @ StorageError::Unsupported(_) => Self::NonRetryable(anyhow::Error::new(error)),
error @ StorageError::RateLimited { .. } => Self::Retry(anyhow::Error::new(error)),
}
}
}
#[async_trait::async_trait]
pub trait StorageClient: Send + Sync + 'static {
async fn open_dataset(&self, name: Option<&str>) -> StorageResult<Arc<dyn Dataset>>;
async fn open_key_value_store(
&self,
name: Option<&str>,
) -> StorageResult<Arc<dyn KeyValueStore>>;
async fn open_request_queue(&self, name: Option<&str>) -> StorageResult<Arc<dyn RequestQueue>>;
async fn purge(&self) -> StorageResult<()>;
}