use std::future::Future;
use base64::prelude::*;
use serde::Serialize;
use std::time::SystemTime;
use crate::{
client::{ClientError, VersaClient, customer_registration::CustomerRegistration},
protocol::{
AssetRegistrationResponse, ClientMetadata, EncryptionKey, EventRegistrationRequest,
EventRegistrationResponse, EventRegistrationSummary, ReceiptRegistrationRequest,
ReceiptRegistrationResponse, ReceiverFilter, ReceiverInstruction, ReceiverQueryRequest,
ReceiverQueryResponse, TransactionHandles,
customer_registration::CustomerRef,
event::{EventType, InitialEventType, UpdateEventType},
webhook::WebhookEvent,
},
};
pub trait VersaSender {
fn register_event(
&self,
event_type: EventType,
handles: Option<TransactionHandles>,
transaction_id: Option<String>,
) -> impl Future<Output = Result<EventRegistrationResponse, ClientError>> + Send;
fn register_initial_event(
&self,
event_type: InitialEventType,
handles: TransactionHandles,
) -> impl Future<Output = Result<EventRegistrationResponse, ClientError>> + Send;
fn register_update_event(
&self,
event_type: UpdateEventType,
transaction_id: String,
) -> impl Future<Output = Result<EventRegistrationResponse, ClientError>> + Send;
#[deprecated(since = "1.6.0", note = "please use `register_event()` instead")]
fn register_receipt(
&self,
handles: TransactionHandles,
transaction_id: Option<String>,
) -> impl Future<Output = Result<ReceiptRegistrationResponse, ClientError>> + Send;
fn encrypt_and_send<T>(
&self,
receiver: ReceiverInstruction,
summary: EventRegistrationSummary,
encryption_key: EncryptionKey,
data: T,
) -> impl Future<Output = Result<(), ClientError>>
where
T: Serialize;
fn register_asset(
&self,
file_data: Vec<u8>,
filename: String,
) -> impl Future<Output = Result<AssetRegistrationResponse, ClientError>> + Send;
fn query_receivers(
&self,
filters: Option<Vec<ReceiverFilter>>,
handles: Option<TransactionHandles>,
) -> impl Future<Output = Result<ReceiverQueryResponse, ClientError>> + Send;
}
pub struct VersaSendingClient {
base_client: VersaClient,
pub schema_version: String,
}
#[cfg(feature = "client_sender")]
impl VersaClient {
pub fn sending_client(self, schema_version: String) -> VersaSendingClient {
VersaSendingClient {
base_client: self,
schema_version,
}
}
}
impl VersaSendingClient {
pub fn client_id(&self) -> String {
self.base_client.client_id()
}
}
#[cfg(feature = "client_sender")]
impl VersaSender for VersaSendingClient {
async fn register_event(
&self,
event_type: EventType,
handles: Option<TransactionHandles>,
transaction_id: Option<String>,
) -> Result<EventRegistrationResponse, ClientError> {
let credential = self.base_client.authorization_header_val();
let schema_version = self.schema_version.clone();
let payload = EventRegistrationRequest {
event_type: Some(event_type),
schema_version,
handles,
transaction_id,
client_metadata: Some(ClientMetadata {
client_string: self.base_client.client_string(),
}),
transaction_event_filter: None,
};
let payload_json = serde_json::to_string(&payload).unwrap();
let url = format!("{}/register", self.base_client.registry_url);
let client = reqwest::Client::new();
let response_result = client
.post(url)
.header("Accept", "application/json")
.header("Authorization", credential)
.header("Content-Type", "application/json")
.body(payload_json)
.send()
.await;
let res = match response_result {
Ok(res) => res,
Err(e) => {
return Err(ClientError::NetworkError(e));
}
};
if res.status().is_success() {
let data: EventRegistrationResponse = match res.json().await {
Ok(val) => val,
Err(e) => {
return Err(ClientError::DeserializationError(e));
}
};
return Ok(data);
} else {
return Err(ClientError::RegistryError(
res.status(),
res.text().await.unwrap_or_default(),
));
}
}
async fn register_initial_event(
&self,
event_type: InitialEventType,
handles: TransactionHandles,
) -> Result<EventRegistrationResponse, ClientError> {
self
.register_event(event_type.into(), Some(handles), None)
.await
}
async fn register_update_event(
&self,
event_type: UpdateEventType,
transaction_id: String,
) -> Result<EventRegistrationResponse, ClientError> {
self
.register_event(event_type.into(), None, Some(transaction_id))
.await
}
async fn register_receipt(
&self,
handles: TransactionHandles,
transaction_id: Option<String>,
) -> Result<ReceiptRegistrationResponse, ClientError> {
let credential = self.base_client.authorization_header_val();
let schema_version = self.schema_version.clone();
let payload = ReceiptRegistrationRequest {
event_type: None,
schema_version,
handles: Some(handles),
transaction_id,
client_metadata: Some(ClientMetadata {
client_string: self.base_client.client_string(),
}),
transaction_event_filter: None,
};
let payload_json = serde_json::to_string(&payload).unwrap();
let url = format!("{}/register", self.base_client.registry_url);
let client = reqwest::Client::new();
let response_result = client
.post(url)
.header("Accept", "application/json")
.header("Authorization", credential)
.header("Content-Type", "application/json")
.body(payload_json)
.send()
.await;
let res = match response_result {
Ok(res) => res,
Err(e) => {
return Err(ClientError::NetworkError(e));
}
};
if res.status().is_success() {
let data: ReceiptRegistrationResponse = match res.json().await {
Ok(val) => val,
Err(e) => {
return Err(ClientError::DeserializationError(e));
}
};
return Ok(data);
} else {
return Err(ClientError::RegistryError(
res.status(),
res.text().await.unwrap_or_default(),
));
}
}
async fn encrypt_and_send<T>(
&self,
receiver: ReceiverInstruction,
summary: EventRegistrationSummary,
encryption_key: EncryptionKey,
data: T,
) -> Result<(), ClientError>
where
T: Serialize,
{
let envelope = crate::encryption::encrypt_envelope(
&data,
&BASE64_STANDARD.decode(encryption_key.0).unwrap(),
);
let EventRegistrationSummary {
mode: _,
event_id: _,
receipt_id,
transaction_id: _,
} = summary;
let data = crate::protocol::ReceiverPayload {
sender_client_id: self.base_client.client_id(),
receipt_id,
envelope,
};
let timestamp = std::time::SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.unwrap()
.as_secs() as i64;
let payload = WebhookEvent {
data,
event_id: Some(receiver.event_id),
event_at: Some(timestamp),
delivery_id: None, delivery_at: Some(timestamp),
event: receiver.event_type.into(),
};
let payload_json = serde_json::to_string(&payload).unwrap();
let byte_body = bytes::Bytes::from(payload_json.clone());
let token = crate::hmac_util::generate_token(byte_body, receiver.secret);
let client = reqwest::Client::new();
let response_result = client
.post(&receiver.address)
.header("Content-Type", "application/json")
.header("X-Request-Signature", token)
.body(payload_json)
.send()
.await;
let res = match response_result {
Ok(res) => res,
Err(e) => {
return Err(ClientError::NetworkError(e));
}
};
if res.status().is_success() {
return Ok(());
} else {
return Err(ClientError::RemoteClientError(
res.status(),
res.text().await.unwrap_or_default(),
));
}
}
async fn register_asset(
&self,
file_data: Vec<u8>,
filename: String,
) -> Result<AssetRegistrationResponse, ClientError> {
let credential = self.base_client.authorization_header_val();
let url = format!("{}/asset", self.base_client.registry_url);
let client = reqwest::Client::new();
let form = reqwest::multipart::Form::new().part(
"file",
reqwest::multipart::Part::bytes(file_data).file_name(filename),
);
let response_result = client
.post(url)
.header("Authorization", credential)
.multipart(form)
.send()
.await;
let res = match response_result {
Ok(res) => res,
Err(e) => {
return Err(ClientError::NetworkError(e));
}
};
if res.status().is_success() {
let data: AssetRegistrationResponse = match res.json().await {
Ok(val) => val,
Err(e) => {
return Err(ClientError::DeserializationError(e));
}
};
return Ok(data);
} else {
return Err(ClientError::RegistryError(
res.status(),
res.text().await.unwrap_or_default(),
));
}
}
async fn query_receivers(
&self,
filters: Option<Vec<ReceiverFilter>>,
handles: Option<TransactionHandles>,
) -> Result<ReceiverQueryResponse, ClientError> {
let credential = self.base_client.authorization_header_val();
let payload = ReceiverQueryRequest { filters, handles };
let payload_json = serde_json::to_string(&payload).unwrap();
let url = format!("{}/receiver/query", self.base_client.registry_url);
let client = reqwest::Client::new();
let response_result = client
.post(url)
.header("Accept", "application/json")
.header("Authorization", credential)
.header("Content-Type", "application/json")
.body(payload_json)
.send()
.await;
let res = match response_result {
Ok(res) => res,
Err(e) => {
return Err(ClientError::NetworkError(e));
}
};
if res.status().is_success() {
let data: ReceiverQueryResponse = match res.json().await {
Ok(val) => val,
Err(e) => {
return Err(ClientError::DeserializationError(e));
}
};
return Ok(data);
} else {
return Err(ClientError::RegistryError(
res.status(),
res.text().await.unwrap_or_default(),
));
}
}
}
impl CustomerRegistration for VersaSendingClient {
async fn register_customer_reference(
&self,
customer_reference: CustomerRef,
) -> Result<(), ClientError> {
self
.base_client
.register_customer_reference(customer_reference)
.await
}
async fn deregister_customer_reference(
&self,
customer_reference: CustomerRef,
) -> Result<(), ClientError> {
self
.base_client
.register_customer_reference(customer_reference)
.await
}
}