1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
#[cfg(feature = "scylla")]
use scylla_driver::transport::errors::QueryError;
use std::sync::{Arc, Mutex};
use thiserror::Error;
use tokio::sync::mpsc::error::SendError as MpscSendError;
use tokio::sync::oneshot::error::RecvError as OneShotRecvError;
#[derive(Debug, Clone, Error)]
pub enum ShardError {
#[error("shard is at capacity, and is unable to evict any records")]
ShardAtCapacity,
#[error("data returned from upstream is immediately expired.")]
DataImmediatelyExpired,
#[error("upstream error: {error:?}")]
UpstreamError { error: UpstreamError },
#[error("the shard is no longer running")]
ShardGone,
}
impl From<OneShotRecvError> for ShardError {
fn from(_error: OneShotRecvError) -> Self {
ShardError::ShardGone
}
}
impl<T> From<MpscSendError<T>> for ShardError {
fn from(_error: MpscSendError<T>) -> Self {
ShardError::ShardGone
}
}
#[derive(Debug, Error, Clone)]
pub enum UpstreamError {
#[error("operation aborted")]
OperationAborted,
#[error("key not found")]
KeyNotFound,
#[error("driver error: {error:?}")]
DriverError { error: Arc<anyhow::Error> },
#[error("serialization error: {error:?}")]
SerializationError { error: Arc<anyhow::Error> },
}
impl UpstreamError {
pub fn is_not_found(&self) -> bool {
matches!(self, UpstreamError::KeyNotFound)
}
pub fn serialization_error<E>(error: E) -> Self
where
E: std::error::Error + Send + 'static,
{
UpstreamError::SerializationError {
error: Arc::new(SyncError::new(error).into()),
}
}
}
impl From<anyhow::Error> for UpstreamError {
fn from(error: anyhow::Error) -> Self {
UpstreamError::DriverError {
error: Arc::new(error),
}
}
}
#[cfg(feature = "cassandra")]
impl From<cassandra_cpp::Error> for UpstreamError {
fn from(error: cassandra_cpp::Error) -> Self {
UpstreamError::DriverError {
error: Arc::new(SyncError::new(error).into()),
}
}
}
#[cfg(feature = "cassandra")]
impl From<cassandra_cpp::Error> for ShardError {
fn from(error: cassandra_cpp::Error) -> Self {
let upstream_error: UpstreamError = error.into();
upstream_error.into()
}
}
#[cfg(feature = "scylla")]
impl From<QueryError> for UpstreamError {
fn from(error: QueryError) -> Self {
UpstreamError::DriverError {
error: Arc::new(error.into()),
}
}
}
impl From<UpstreamError> for ShardError {
fn from(error: UpstreamError) -> Self {
ShardError::UpstreamError { error }
}
}
pub struct SyncError<E> {
inner: Mutex<E>,
}
impl<E: std::error::Error + Send + 'static> SyncError<E> {
pub fn new(err: E) -> Self {
SyncError {
inner: Mutex::new(err),
}
}
}
impl<T> std::fmt::Display for SyncError<T>
where
T: std::fmt::Display,
{
fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
self.inner.lock().unwrap().fmt(f)
}
}
impl<T> std::fmt::Debug for SyncError<T>
where
T: std::fmt::Debug,
{
fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
(*self.inner.lock().unwrap()).fmt(f)
}
}
impl<E: std::error::Error + Send + 'static> std::error::Error for SyncError<E> {}