use std::sync::Arc;
use crate::{
ExportVid, OwnedVid, PrivateVid, RelationshipStatus,
crypto::CryptoError,
definitions::{Digest, ReceivedTspMessage, TSPStream, VerifiedVid},
error::Error,
store::{Aliases, SecureStore, WebvhUpdateKeys},
};
use bytes::BytesMut;
use futures::StreamExt;
use tracing::debug;
use url::Url;
#[derive(Default, Clone)]
pub struct AsyncSecureStore {
inner: SecureStore,
}
impl AsyncSecureStore {
pub fn new() -> Self {
Default::default()
}
pub fn export(&self) -> Result<(Vec<ExportVid>, Aliases, WebvhUpdateKeys), Error> {
self.inner.export()
}
pub fn as_store(&self) -> &SecureStore {
&self.inner
}
pub fn import(
&self,
vids: Vec<ExportVid>,
aliases: Aliases,
keys: WebvhUpdateKeys,
) -> Result<(), Error> {
self.inner.import(vids, aliases, keys)
}
pub fn get_relation_status_for_vid_pair(
&self,
local_vid: &str,
remote_vid: &str,
) -> Result<RelationshipStatus, Error> {
self.inner
.relation_status_for_vid_pair(local_vid, remote_vid)
}
pub fn set_relation_and_status_for_vid(
&self,
vid: &str,
status: RelationshipStatus,
relation_vid: &str,
) -> Result<(), Error> {
self.inner
.set_relation_and_status_for_vid(vid, status, relation_vid)
}
pub fn set_route_for_vid(&self, vid: &str, route: &[&str]) -> Result<(), Error> {
self.inner.set_route_for_vid(vid, route)
}
pub fn set_parent_for_vid(&self, vid: &str, parent: Option<&str>) -> Result<(), Error> {
self.inner.set_parent_for_vid(vid, parent)
}
pub fn list_vids(&self) -> Result<Vec<String>, Error> {
self.inner.list_vids()
}
pub fn add_private_vid(
&self,
private_vid: impl PrivateVid + Clone + 'static,
metadata: Option<serde_json::Value>,
) -> Result<(), Error> {
self.inner.add_private_vid(private_vid, metadata)
}
pub fn forget_vid(&self, vid: &str) -> Result<(), Error> {
self.inner.forget_vid(vid)
}
pub fn add_verified_vid(
&self,
verified_vid: impl VerifiedVid + 'static,
metadata: Option<serde_json::Value>,
) -> Result<(), Error> {
self.inner.add_verified_vid(verified_vid, metadata)
}
pub fn has_private_vid(&self, vid: &str) -> Result<bool, Error> {
self.inner.has_private_vid(vid)
}
pub fn has_verified_vid(&self, vid: &str) -> Result<bool, Error> {
self.inner.has_verified_vid(vid)
}
pub fn get_verified_vid(&self, vid: &str) -> Result<Arc<dyn VerifiedVid>, Error> {
self.inner.get_verified_vid(vid)
}
pub async fn verify_vid(&self, vid: &str, alias: Option<String>) -> Result<(), Error> {
let (verified_vid, metadata) = crate::vid::verify_vid(vid).await?;
self.inner.add_verified_vid(verified_vid, metadata)?;
if let Some(alias) = alias {
self.set_alias(alias, vid.to_owned())?;
}
Ok(())
}
pub fn resolve_alias(&self, alias: &str) -> Result<Option<String>, Error> {
self.inner.resolve_alias(alias)
}
pub fn try_resolve_alias(&self, alias: &str) -> Result<String, Error> {
self.inner.try_resolve_alias(alias)
}
pub fn set_alias(&self, alias: String, did: String) -> Result<(), Error> {
self.inner.set_alias(alias, did)
}
pub fn add_secret_key(&self, kid: String, secret_key: Vec<u8>) -> Result<(), Error> {
self.inner.add_secret_key(kid, secret_key)
}
pub fn get_secret_key(&self, kid: &str) -> Result<Option<Vec<u8>>, Error> {
self.inner.get_secret_key(kid)
}
pub fn seal_message(
&self,
sender: &str,
receiver: &str,
nonconfidential_data: Option<&[u8]>,
message: &[u8],
) -> Result<(Url, Vec<u8>), Error> {
self.inner
.seal_message(sender, receiver, nonconfidential_data, message)
}
pub async fn send(
&self,
sender: &str,
receiver: &str,
nonconfidential_data: Option<&[u8]>,
message: &[u8],
) -> Result<(), Error> {
match self.inner.relation_status_for_vid_pair(sender, receiver) {
Ok(relation) => {
if matches!(relation, RelationshipStatus::Unrelated) {
self.send_relationship_request(sender, receiver, None)
.await?
}
}
Err(Error::Relationship(_)) => {
self.send_relationship_request(sender, receiver, None)
.await?
}
Err(e) => return Err(e),
};
let (endpoint, message) =
self.inner
.seal_message(sender, receiver, nonconfidential_data, message)?;
tracing::info!("sending message to {endpoint}");
crate::transport::send_message(&endpoint, &message).await?;
Ok(())
}
pub fn make_relationship_request(
&self,
sender: &str,
receiver: &str,
route: Option<&[&str]>,
) -> Result<(Url, Vec<u8>), Error> {
self.inner
.make_relationship_request(sender, receiver, route)
}
pub async fn send_relationship_request(
&self,
sender: &str,
receiver: &str,
route: Option<&[&str]>,
) -> Result<(), Error> {
let (endpoint, message) = self
.inner
.make_relationship_request(sender, receiver, route)?;
tracing::info!("sending message to {endpoint}");
crate::transport::send_message(&endpoint, &message).await?;
Ok(())
}
pub fn make_relationship_accept(
&self,
sender: &str,
receiver: &str,
thread_id: Digest,
route: Option<&[&str]>,
) -> Result<(Url, Vec<u8>), Error> {
self.inner
.make_relationship_accept(sender, receiver, thread_id, route)
}
pub async fn send_relationship_accept(
&self,
sender: &str,
receiver: &str,
thread_id: Digest,
route: Option<&[&str]>,
) -> Result<(), Error> {
let (endpoint, message) =
self.make_relationship_accept(sender, receiver, thread_id, route)?;
tracing::info!("sending message to {endpoint}");
crate::transport::send_message(&endpoint, &message).await?;
Ok(())
}
pub fn make_relationship_cancel(
&self,
sender: &str,
receiver: &str,
) -> Result<(Url, Vec<u8>), Error> {
self.inner.make_relationship_cancel(sender, receiver)
}
pub async fn send_relationship_cancel(
&self,
sender: &str,
receiver: &str,
) -> Result<(), Error> {
let (endpoint, message) = self.inner.make_relationship_cancel(sender, receiver)?;
tracing::info!("sending message to {endpoint}");
crate::transport::send_message(&endpoint, &message).await?;
Ok(())
}
pub fn make_new_identifier_notice(
&self,
sender: &str,
receiver: &str,
sender_new_vid: &str,
) -> Result<(Url, Vec<u8>), Error> {
self.inner
.make_new_identifier_notice(sender, receiver, sender_new_vid)
}
pub async fn send_new_identifier_notice(
&self,
sender: &str,
receiver: &str,
sender_new_vid: &str,
) -> Result<(), Error> {
let (endpoint, message) =
self.inner
.make_new_identifier_notice(sender, receiver, sender_new_vid)?;
tracing::info!("sending message to {endpoint}");
crate::transport::send_message(&endpoint, &message).await?;
Ok(())
}
pub fn make_relationship_referral(
&self,
sender: &str,
receiver: &str,
referred_vid: &str,
) -> Result<(Url, Vec<u8>), Error> {
self.inner
.make_relationship_referral(sender, receiver, referred_vid)
}
pub async fn send_relationship_referral(
&self,
sender: &str,
receiver: &str,
referred_vid: &str,
) -> Result<(), Error> {
let (endpoint, message) =
self.inner
.make_relationship_referral(sender, receiver, referred_vid)?;
tracing::info!("sending message to {endpoint}");
crate::transport::send_message(&endpoint, &message).await?;
Ok(())
}
pub fn make_nested_relationship_request(
&self,
parent_sender: &str,
receiver: &str,
) -> Result<((Url, Vec<u8>), OwnedVid), Error> {
self.inner
.make_nested_relationship_request(parent_sender, receiver)
}
pub async fn send_nested_relationship_request(
&self,
parent_sender: &str,
receiver: &str,
) -> Result<OwnedVid, Error> {
let ((endpoint, message), vid) = self
.inner
.make_nested_relationship_request(parent_sender, receiver)?;
tracing::info!("sending message to {endpoint}");
crate::transport::send_message(&endpoint, &message).await?;
Ok(vid)
}
pub fn make_nested_relationship_accept(
&self,
parent_sender: &str,
nested_receiver: &str,
thread_id: Digest,
) -> Result<((Url, Vec<u8>), OwnedVid), Error> {
self.inner
.make_nested_relationship_accept(parent_sender, nested_receiver, thread_id)
}
pub async fn send_nested_relationship_accept(
&self,
parent_sender: &str,
nested_receiver: &str,
thread_id: Digest,
) -> Result<OwnedVid, Error> {
let ((endpoint, message), vid) =
self.make_nested_relationship_accept(parent_sender, nested_receiver, thread_id)?;
tracing::info!("sending message to {endpoint}");
crate::transport::send_message(&endpoint, &message).await?;
Ok(vid)
}
pub fn make_next_routed_message(
&self,
next_hop: &str,
path: Vec<impl AsRef<[u8]>>,
opaque_message: &[u8],
) -> Result<(Url, Vec<u8>), Error> {
self.inner.forward_routed_message(
next_hop,
path.iter().map(|x| x.as_ref()).collect(),
opaque_message,
)
}
pub async fn forward_routed_message(
&self,
next_hop: &str,
path: Vec<impl AsRef<[u8]>>,
opaque_message: &[u8],
) -> Result<Url, Error> {
let (transport, message) = self.make_next_routed_message(next_hop, path, opaque_message)?;
crate::transport::send_message(&transport, &message).await?;
Ok(transport)
}
pub fn open_message<'a>(
&self,
message: &'a mut [u8],
) -> Result<ReceivedTspMessage<&'a [u8]>, Error> {
self.inner.open_message(message)
}
pub async fn receive(&self, vid: &str) -> Result<TSPStream<ReceivedTspMessage, Error>, Error> {
let receiver = self.inner.get_private_vid(vid)?;
let mut transport = receiver.endpoint().clone();
let path = transport
.path()
.replace("[vid_placeholder]", &self.inner.try_resolve_alias(vid)?);
transport.set_path(&path);
tracing::trace!("Listening for {vid} on {transport}");
let messages = crate::transport::receive_messages(&transport).await?;
let db = self.inner.clone();
let self_clone = self.clone();
Ok(Box::pin(messages.then(move |message| {
let db_inner = db.clone();
let self_inner = self_clone.clone();
async move {
let mut message = message?;
if tracing::enabled!(tracing::Level::TRACE) {
println!(
"CESR-encoded message: {}",
crate::cesr::color_format(&message)?
);
}
match db_inner.open_message(&mut message) {
Err(Error::UnverifiedSource(unknown_vid, _)) => {
debug!("Verifying VID: {}", unknown_vid);
self_inner.verify_vid(&unknown_vid, None).await?;
db_inner.open_message(&mut message)
}
Err(Error::Crypto(CryptoError::Verify(vid, _))) => {
debug!("Re-verifying VID: {}", vid);
self_inner.verify_vid(&vid, None).await?;
db_inner.open_message(&mut message)
}
maybe_message => maybe_message,
}
.map(|msg| msg.into_owned())
.map_err(|e| {
tracing::error!("{}", e);
e
})
}
})))
}
pub async fn send_anycast(
&self,
sender: &str,
receivers: impl IntoIterator<Item = impl AsRef<str>>,
nonconfidential_message: &[u8],
) -> Result<(), Error> {
let message = self.inner.sign_anycast(sender, nonconfidential_message)?;
for vid in receivers {
let receiver = self.inner.get_verified_vid(vid.as_ref())?;
crate::transport::send_message(receiver.endpoint(), &message).await?;
}
Ok(())
}
pub async fn verify_and_open(
&self,
vid: &str,
mut payload: BytesMut,
) -> Result<ReceivedTspMessage, Error> {
self.verify_vid(vid, None).await?;
Ok(self.inner.open_message(&mut payload)?.into_owned())
}
}