mod cmd;
mod data_exchange;
mod errors;
mod query;
mod register;
pub use self::{
cmd::DataCmd,
data_exchange::{ChunkDataExchange, DataExchange, RegisterDataExchange, StorageLevel},
errors::{Error, Result},
query::DataQuery,
register::{RegisterCmd, RegisterRead, RegisterWrite},
};
use crate::types::{
register::{Entry, EntryHash, Permissions, Policy, Register},
Chunk, ChunkAddress, DataAddress, PublicKey,
};
use crate::{
messaging::{data::Error as ErrorMessage, MessageId},
types::utils,
};
use bytes::Bytes;
use serde::{Deserialize, Serialize};
use std::{collections::BTreeSet, convert::TryFrom};
use xor_name::XorName;
pub type OperationId = String;
pub fn operation_id(address: &ChunkAddress) -> Result<OperationId> {
utils::encode(address).map_err(|_| Error::NoOperationId)
}
#[allow(clippy::large_enum_variant)]
#[derive(Debug, Eq, PartialEq, Clone, Serialize, Deserialize)]
pub struct ServiceError {
pub reason: Option<Error>,
pub source_message: Option<Bytes>,
}
#[allow(clippy::large_enum_variant)]
#[derive(Debug, Eq, PartialEq, Clone, Serialize, Deserialize)]
pub enum ServiceMsg {
Cmd(DataCmd),
Query(DataQuery),
QueryResponse {
response: QueryResponse,
correlation_id: MessageId,
},
CmdError {
error: CmdError,
correlation_id: MessageId,
},
ServiceError(ServiceError),
}
impl ServiceMsg {
pub fn dst_address(&self) -> Option<XorName> {
match self {
Self::Cmd(cmd) => Some(cmd.dst_name()),
Self::Query(query) => Some(query.dst_name()),
_ => None,
}
}
}
#[derive(Debug, Hash, Eq, PartialEq, Clone, Serialize, Deserialize)]
pub enum CmdError {
Data(Error), }
#[allow(clippy::large_enum_variant, clippy::type_complexity)]
#[derive(Eq, PartialEq, Clone, Serialize, Deserialize, Debug)]
pub enum QueryResponse {
GetChunk(Result<Chunk>),
GetRegister((Result<Register>, OperationId)),
GetRegisterOwner((Result<PublicKey>, OperationId)),
ReadRegister((Result<BTreeSet<(EntryHash, Entry)>>, OperationId)),
GetRegisterPolicy((Result<Policy>, OperationId)),
GetRegisterUserPermissions((Result<Permissions>, OperationId)),
}
impl QueryResponse {
pub fn is_success(&self) -> bool {
use QueryResponse::*;
match self {
GetChunk(result) => result.is_ok(),
GetRegister((result, _op_id)) => result.is_ok(),
GetRegisterOwner((result, _op_id)) => result.is_ok(),
ReadRegister((result, _op_id)) => result.is_ok(),
GetRegisterPolicy((result, _op_id)) => result.is_ok(),
GetRegisterUserPermissions((result, _op_id)) => result.is_ok(),
}
}
pub fn failed_with_data_not_found(&self) -> bool {
use QueryResponse::*;
match self {
GetChunk(result) => match result {
Ok(_) => false,
Err(error) => matches!(*error, ErrorMessage::ChunkNotFound(_)),
},
GetRegister((result, _op_id)) => match result {
Ok(_) => false,
Err(error) => matches!(*error, ErrorMessage::DataNotFound(_)),
},
GetRegisterOwner((result, _op_id)) => match result {
Ok(_) => false,
Err(error) => matches!(*error, ErrorMessage::DataNotFound(_)),
},
ReadRegister((result, _op_id)) => match result {
Ok(_) => false,
Err(error) => matches!(*error, ErrorMessage::DataNotFound(_)),
},
GetRegisterPolicy((result, _op_id)) => match result {
Ok(_) => false,
Err(error) => matches!(*error, ErrorMessage::DataNotFound(_)),
},
GetRegisterUserPermissions((result, _op_id)) => match result {
Ok(_) => false,
Err(error) => matches!(*error, ErrorMessage::DataNotFound(_)),
},
}
}
pub fn operation_id(&self) -> Result<OperationId> {
use QueryResponse::*;
match self {
GetChunk(result) => match result {
Ok(chunk) => operation_id(chunk.address()),
Err(ErrorMessage::ChunkNotFound(name)) => operation_id(&ChunkAddress(*name)),
Err(ErrorMessage::DataNotFound(DataAddress::Bytes(address))) => {
operation_id(&ChunkAddress(*address.name()))
}
Err(ErrorMessage::DataNotFound(another_address)) => {
error!(
"{:?} address returned when we were expecting a ChunkAddress",
another_address
);
Err(Error::NoOperationId)
}
Err(another_error) => {
error!("Could not form operation id: {:?}", another_error);
Err(Error::InvalidQueryResponseErrorForOperationId)
}
},
GetRegister((_, operation_id))
| GetRegisterOwner((_, operation_id))
| ReadRegister((_, operation_id))
| GetRegisterPolicy((_, operation_id))
| GetRegisterUserPermissions((_, operation_id)) => Ok(operation_id.clone()),
}
}
}
#[derive(Debug, PartialEq)]
#[allow(clippy::large_enum_variant)]
pub enum TryFromError {
WrongType,
Response(Error),
}
macro_rules! try_from {
($ok_type:ty, $($variant:ident),*) => {
impl TryFrom<QueryResponse> for $ok_type {
type Error = TryFromError;
fn try_from(response: QueryResponse) -> std::result::Result<Self, Self::Error> {
match response {
$(
QueryResponse::$variant((Ok(data), _op_id)) => Ok(data),
QueryResponse::$variant((Err(error), _op_id)) => Err(TryFromError::Response(error)),
)*
_ => Err(TryFromError::WrongType),
}
}
}
};
}
impl TryFrom<QueryResponse> for Chunk {
type Error = TryFromError;
fn try_from(response: QueryResponse) -> std::result::Result<Self, Self::Error> {
match response {
QueryResponse::GetChunk(Ok(data)) => Ok(data),
QueryResponse::GetChunk(Err(error)) => Err(TryFromError::Response(error)),
_ => Err(TryFromError::WrongType),
}
}
}
try_from!(Register, GetRegister);
try_from!(PublicKey, GetRegisterOwner);
try_from!(BTreeSet<(EntryHash, Entry)>, ReadRegister);
try_from!(Policy, GetRegisterPolicy);
try_from!(Permissions, GetRegisterUserPermissions);
#[cfg(test)]
mod tests {
use super::*;
use crate::types::{utils::random_bytes, Chunk, Keypair};
use bytes::Bytes;
use eyre::{eyre, Result};
use std::convert::{TryFrom, TryInto};
fn gen_keypairs() -> Vec<Keypair> {
let mut rng = rand::thread_rng();
let bls_secret_key = bls::SecretKeySet::random(1, &mut rng);
vec![
Keypair::new_ed25519(&mut rng),
Keypair::new_bls_share(
0,
bls_secret_key.secret_key_share(0),
bls_secret_key.public_keys(),
),
]
}
pub(crate) fn gen_keys() -> Vec<PublicKey> {
gen_keypairs().iter().map(PublicKey::from).collect()
}
#[test]
fn debug_format_functional() -> Result<()> {
if let Some(key) = gen_keys().first() {
let errored_response = QueryResponse::GetRegister((
Err(Error::AccessDenied(*key)),
"some_op_id".to_string(),
));
assert!(format!("{:?}", errored_response).contains("GetRegister((Err(AccessDenied("));
Ok(())
} else {
Err(eyre!("Could not generate public key"))
}
}
#[test]
fn try_from() -> Result<()> {
use QueryResponse::*;
let key = match gen_keys().first() {
Some(key) => *key,
None => return Err(eyre!("Could not generate public key")),
};
let i_data = Chunk::new(Bytes::from(vec![1, 3, 1, 4]));
let e = Error::AccessDenied(key);
assert_eq!(
i_data,
GetChunk(Ok(i_data.clone()))
.try_into()
.map_err(|_| eyre!("Mismatched types".to_string()))?
);
assert_eq!(
Err(TryFromError::Response(e.clone())),
Chunk::try_from(GetChunk(Err(e)))
);
Ok(())
}
#[test]
fn wire_msg_payload() -> Result<()> {
use crate::messaging::data::DataCmd;
use crate::messaging::data::ServiceMsg;
use crate::messaging::WireMsg;
let chunks = (0..10).map(|_| Chunk::new(random_bytes(3072)));
for chunk in chunks {
let (original_msg, serialised_cmd) = {
let msg = ServiceMsg::Cmd(DataCmd::StoreChunk(chunk));
let bytes = WireMsg::serialize_msg_payload(&msg)?;
(msg, bytes)
};
let deserialized_msg: ServiceMsg =
rmp_serde::from_slice(&serialised_cmd).map_err(|err| {
crate::messaging::Error::FailedToParse(format!(
"Data message payload as Msgpack: {}",
err
))
})?;
assert_eq!(original_msg, deserialized_msg);
}
Ok(())
}
}