use protobuf::Message as M;
use protobuf::RepeatedField;
use messages::events::Event;
use messages::events::Event_Attribute;
use messages::state_context::*;
use messages::validator::Message_MessageType;
use messaging::stream::MessageSender;
use messaging::zmq_stream::ZmqMessageSender;
use processor::handler::{ContextError, TransactionContext};
use super::generate_correlation_id;
#[derive(Clone)]
pub struct ZmqTransactionContext {
context_id: String,
sender: ZmqMessageSender,
}
impl ZmqTransactionContext {
pub fn new(context_id: &str, sender: ZmqMessageSender) -> Self {
ZmqTransactionContext {
context_id: String::from(context_id),
sender,
}
}
}
impl TransactionContext for ZmqTransactionContext {
fn get_state_entries(
&self,
addresses: &[String],
) -> Result<Vec<(String, Vec<u8>)>, ContextError> {
let mut request = TpStateGetRequest::new();
request.set_context_id(self.context_id.clone());
request.set_addresses(RepeatedField::from_vec(addresses.to_vec()));
let serialized = request.write_to_bytes()?;
let x: &[u8] = &serialized;
let mut future = self.sender.send(
Message_MessageType::TP_STATE_GET_REQUEST,
&generate_correlation_id(),
x,
)?;
let response: TpStateGetResponse = protobuf::parse_from_bytes(future.get()?.get_content())?;
match response.get_status() {
TpStateGetResponse_Status::OK => {
let mut entries = Vec::new();
for entry in response.get_entries() {
match entry.get_data().len() {
0 => continue,
_ => entries
.push((entry.get_address().to_string(), Vec::from(entry.get_data()))),
}
}
Ok(entries)
}
TpStateGetResponse_Status::AUTHORIZATION_ERROR => {
Err(ContextError::AuthorizationError(format!(
"Tried to get unauthorized addresses: {:?}",
addresses
)))
}
TpStateGetResponse_Status::STATUS_UNSET => Err(ContextError::ResponseAttributeError(
String::from("Status was not set for TpStateGetResponse"),
)),
}
}
fn set_state_entries(&self, entries: Vec<(String, Vec<u8>)>) -> Result<(), ContextError> {
let state_entries: Vec<TpStateEntry> = entries
.into_iter()
.map(|(address, payload)| {
let mut entry = TpStateEntry::new();
entry.set_address(address);
entry.set_data(payload);
entry
})
.collect();
let mut request = TpStateSetRequest::new();
request.set_context_id(self.context_id.clone());
request.set_entries(RepeatedField::from_vec(state_entries.to_vec()));
let serialized = request.write_to_bytes()?;
let x: &[u8] = &serialized;
let mut future = self.sender.send(
Message_MessageType::TP_STATE_SET_REQUEST,
&generate_correlation_id(),
x,
)?;
let response: TpStateSetResponse = protobuf::parse_from_bytes(future.get()?.get_content())?;
match response.get_status() {
TpStateSetResponse_Status::OK => Ok(()),
TpStateSetResponse_Status::AUTHORIZATION_ERROR => {
Err(ContextError::AuthorizationError(format!(
"Tried to set unauthorized addresses: {:?}",
state_entries
)))
}
TpStateSetResponse_Status::STATUS_UNSET => Err(ContextError::ResponseAttributeError(
String::from("Status was not set for TpStateSetResponse"),
)),
}
}
fn delete_state_entries(&self, addresses: &[String]) -> Result<Vec<String>, ContextError> {
let mut request = TpStateDeleteRequest::new();
request.set_context_id(self.context_id.clone());
request.set_addresses(RepeatedField::from_slice(addresses));
let serialized = request.write_to_bytes()?;
let x: &[u8] = &serialized;
let mut future = self.sender.send(
Message_MessageType::TP_STATE_DELETE_REQUEST,
&generate_correlation_id(),
x,
)?;
let response: TpStateDeleteResponse =
protobuf::parse_from_bytes(future.get()?.get_content())?;
match response.get_status() {
TpStateDeleteResponse_Status::OK => Ok(Vec::from(response.get_addresses())),
TpStateDeleteResponse_Status::AUTHORIZATION_ERROR => {
Err(ContextError::AuthorizationError(format!(
"Tried to delete unauthorized addresses: {:?}",
addresses
)))
}
TpStateDeleteResponse_Status::STATUS_UNSET => {
Err(ContextError::ResponseAttributeError(String::from(
"Status was not set for TpStateDeleteResponse",
)))
}
}
}
fn add_receipt_data(&self, data: &[u8]) -> Result<(), ContextError> {
let mut request = TpReceiptAddDataRequest::new();
request.set_context_id(self.context_id.clone());
request.set_data(Vec::from(data));
let serialized = request.write_to_bytes()?;
let x: &[u8] = &serialized;
let mut future = self.sender.send(
Message_MessageType::TP_RECEIPT_ADD_DATA_REQUEST,
&generate_correlation_id(),
x,
)?;
let response: TpReceiptAddDataResponse =
protobuf::parse_from_bytes(future.get()?.get_content())?;
match response.get_status() {
TpReceiptAddDataResponse_Status::OK => Ok(()),
TpReceiptAddDataResponse_Status::ERROR => Err(ContextError::TransactionReceiptError(
format!("Failed to add receipt data {:?}", data),
)),
TpReceiptAddDataResponse_Status::STATUS_UNSET => {
Err(ContextError::ResponseAttributeError(String::from(
"Status was not set for TpReceiptAddDataResponse",
)))
}
}
}
fn add_event(
&self,
event_type: String,
attributes: Vec<(String, String)>,
data: &[u8],
) -> Result<(), ContextError> {
let mut event = Event::new();
event.set_event_type(event_type);
let mut attributes_vec = Vec::new();
for (key, value) in attributes {
let mut attribute = Event_Attribute::new();
attribute.set_key(key);
attribute.set_value(value);
attributes_vec.push(attribute);
}
event.set_attributes(RepeatedField::from_vec(attributes_vec));
event.set_data(Vec::from(data));
let mut request = TpEventAddRequest::new();
request.set_context_id(self.context_id.clone());
request.set_event(event.clone());
let serialized = request.write_to_bytes()?;
let x: &[u8] = &serialized;
let mut future = self.sender.send(
Message_MessageType::TP_EVENT_ADD_REQUEST,
&generate_correlation_id(),
x,
)?;
let response: TpEventAddResponse = protobuf::parse_from_bytes(future.get()?.get_content())?;
match response.get_status() {
TpEventAddResponse_Status::OK => Ok(()),
TpEventAddResponse_Status::ERROR => Err(ContextError::TransactionReceiptError(
format!("Failed to add event {:?}", event),
)),
TpEventAddResponse_Status::STATUS_UNSET => Err(ContextError::ResponseAttributeError(
String::from("Status was not set for TpEventAddRespons"),
)),
}
}
}