Skip to content
This repository has been archived by the owner on Oct 31, 2024. It is now read-only.

Commit

Permalink
feat: switch tce-lib action to spawn tasks
Browse files Browse the repository at this point in the history
Signed-off-by: Simon Paitrault <simon.paitrault@gmail.com>
  • Loading branch information
Freyskeyd committed Mar 28, 2024
1 parent b8cd730 commit d695440
Show file tree
Hide file tree
Showing 3 changed files with 131 additions and 115 deletions.
140 changes: 72 additions & 68 deletions crates/topos-tce/src/app_context/api.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
use crate::AppContext;
use std::collections::HashMap;
use tokio::spawn;
use topos_core::uci::{Certificate, SubnetId};
use topos_metrics::CERTIFICATE_DELIVERY_LATENCY;
use topos_tce_api::RuntimeError;
Expand All @@ -20,79 +21,82 @@ impl AppContext {
self.delivery_latency
.insert(certificate.id, CERTIFICATE_DELIVERY_LATENCY.start_timer());

_ = match self
.validator_store
.insert_pending_certificate(&certificate)
.await
{
Ok(Some(pending_id)) => {
let certificate_id = certificate.id;
debug!(
"Certificate {} from subnet {} has been inserted into pending pool",
certificate_id, certificate.source_subnet_id
);
let validator_store = self.validator_store.clone();
let double_echo = self.tce_cli.get_double_echo_channel();

if self
.tce_cli
.get_double_echo_channel()
.send(DoubleEchoCommand::Broadcast {
need_gossip: true,
cert: *certificate,
pending_id,
})
.await
.is_err()
{
error!(
"Unable to send DoubleEchoCommand::Broadcast command to double \
echo for {}",
certificate_id
spawn(async move {
_ = match validator_store
.insert_pending_certificate(&certificate)
.await
{
Ok(Some(pending_id)) => {
let certificate_id = certificate.id;
debug!(
"Certificate {} from subnet {} has been inserted into pending pool",
certificate_id, certificate.source_subnet_id
);

sender.send(Err(RuntimeError::CommunicationError(
"Unable to send DoubleEchoCommand::Broadcast command to double \
echo"
.to_string(),
)))
} else {
sender.send(Ok(PendingResult::InPending(pending_id)))
if double_echo
.send(DoubleEchoCommand::Broadcast {
need_gossip: true,
cert: *certificate,
pending_id,
})
.await
.is_err()
{
error!(
"Unable to send DoubleEchoCommand::Broadcast command to \
double echo for {}",
certificate_id
);

sender.send(Err(RuntimeError::CommunicationError(
"Unable to send DoubleEchoCommand::Broadcast command to \
double echo"
.to_string(),
)))
} else {
sender.send(Ok(PendingResult::InPending(pending_id)))
}
}
}
Ok(None) => {
debug!(
"Certificate {} from subnet {} has been inserted into precedence pool \
waiting for {}",
certificate.id, certificate.source_subnet_id, certificate.prev_id
);
sender.send(Ok(PendingResult::AwaitPrecedence))
}
Err(StorageError::InternalStorage(
InternalStorageError::CertificateAlreadyPending,
)) => {
debug!(
"Certificate {} has already been added to the pending pool, skipping",
certificate.id
);
sender.send(Ok(PendingResult::AlreadyPending))
}
Err(StorageError::InternalStorage(
InternalStorageError::CertificateAlreadyExists,
)) => {
debug!(
"Certificate {} has already been delivered, skipping",
certificate.id
);
sender.send(Ok(PendingResult::AlreadyDelivered))
}
Err(error) => {
error!(
"Unable to insert pending certificate {}: {}",
certificate.id, error
);
Ok(None) => {
debug!(
"Certificate {} from subnet {} has been inserted into precedence \
pool waiting for {}",
certificate.id, certificate.source_subnet_id, certificate.prev_id
);
sender.send(Ok(PendingResult::AwaitPrecedence))
}
Err(StorageError::InternalStorage(
InternalStorageError::CertificateAlreadyPending,
)) => {
debug!(
"Certificate {} has already been added to the pending pool, \
skipping",
certificate.id
);
sender.send(Ok(PendingResult::AlreadyPending))
}
Err(StorageError::InternalStorage(
InternalStorageError::CertificateAlreadyExists,
)) => {
debug!(
"Certificate {} has already been delivered, skipping",
certificate.id
);
sender.send(Ok(PendingResult::AlreadyDelivered))
}
Err(error) => {
error!(
"Unable to insert pending certificate {}: {}",
certificate.id, error
);

sender.send(Err(error.into()))
}
};
sender.send(Err(error.into()))
}
};
});
}

ApiEvent::GetSourceHead { subnet_id, sender } => {
Expand Down
88 changes: 46 additions & 42 deletions crates/topos-tce/src/app_context/protocol.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
use tokio::spawn;
use topos_core::api::grpc::tce::v1::{double_echo_request, DoubleEchoRequest, Echo, Gossip, Ready};
use topos_tce_broadcast::event::ProtocolEvents;
use tracing::{error, info, warn};
Expand All @@ -14,65 +15,68 @@ impl AppContext {
ProtocolEvents::Gossip { cert } => {
let cert_id = cert.id;

let request = DoubleEchoRequest {
request: Some(double_echo_request::Request::Gossip(Gossip {
certificate: Some(cert.into()),
})),
};
let network_client = self.network_client.clone();
spawn(async move {
let request = DoubleEchoRequest {
request: Some(double_echo_request::Request::Gossip(Gossip {
certificate: Some(cert.into()),
})),
};

info!("Sending Gossip for certificate {}", cert_id);
if let Err(e) = self
.network_client
.publish(topos_p2p::TOPOS_GOSSIP, request)
.await
{
error!("Unable to send Gossip: {e}");
}
info!("Sending Gossip for certificate {}", cert_id);
if let Err(e) = network_client
.publish(topos_p2p::TOPOS_GOSSIP, request)
.await
{
error!("Unable to send Gossip: {e}");
}
});
}

ProtocolEvents::Echo {
certificate_id,
signature,
validator_id,
} if self.is_validator => {
// Send echo message
let request = DoubleEchoRequest {
request: Some(double_echo_request::Request::Echo(Echo {
certificate_id: Some(certificate_id.into()),
signature: Some(signature.into()),
validator_id: Some(validator_id.into()),
})),
};
let network_client = self.network_client.clone();
spawn(async move {
// Send echo message
let request = DoubleEchoRequest {
request: Some(double_echo_request::Request::Echo(Echo {
certificate_id: Some(certificate_id.into()),
signature: Some(signature.into()),
validator_id: Some(validator_id.into()),
})),
};

if let Err(e) = self
.network_client
.publish(topos_p2p::TOPOS_ECHO, request)
.await
{
error!("Unable to send Echo: {e}");
}
if let Err(e) = network_client.publish(topos_p2p::TOPOS_ECHO, request).await {
error!("Unable to send Echo: {e}");
}
});
}

ProtocolEvents::Ready {
certificate_id,
signature,
validator_id,
} if self.is_validator => {
let request = DoubleEchoRequest {
request: Some(double_echo_request::Request::Ready(Ready {
certificate_id: Some(certificate_id.into()),
signature: Some(signature.into()),
validator_id: Some(validator_id.into()),
})),
};
let network_client = self.network_client.clone();
spawn(async move {
let request = DoubleEchoRequest {
request: Some(double_echo_request::Request::Ready(Ready {
certificate_id: Some(certificate_id.into()),
signature: Some(signature.into()),
validator_id: Some(validator_id.into()),
})),
};

if let Err(e) = self
.network_client
.publish(topos_p2p::TOPOS_READY, request)
.await
{
error!("Unable to send Ready: {e}");
}
if let Err(e) = network_client
.publish(topos_p2p::TOPOS_READY, request)
.await
{
error!("Unable to send Ready: {e}");
}
});
}
ProtocolEvents::BroadcastFailed { certificate_id } => {
warn!("Broadcast failed for certificate {certificate_id}")
Expand Down
18 changes: 13 additions & 5 deletions crates/topos-tce/src/tests/mod.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
use libp2p::PeerId;
use rstest::{fixture, rstest};
use std::{collections::HashSet, future::IntoFuture, sync::Arc};
use std::{collections::HashSet, future::IntoFuture, sync::Arc, time::Duration};
use tokio_stream::Stream;
use topos_tce_api::RuntimeEvent;
use topos_tce_broadcast::event::ProtocolEvents;
Expand Down Expand Up @@ -41,8 +41,8 @@ async fn non_validator_publish_gossip(
.await;

assert!(matches!(
p2p_receiver.try_recv(),
Ok(topos_p2p::Command::Gossip { topic, .. }) if topic == "topos_gossip"
p2p_receiver.recv().await,
Some(topos_p2p::Command::Gossip { topic, .. }) if topic == "topos_gossip"
));
}

Expand All @@ -64,7 +64,11 @@ async fn non_validator_do_not_publish_echo(
})
.await;

assert!(p2p_receiver.try_recv().is_err(),);
assert!(
tokio::time::timeout(Duration::from_millis(10), p2p_receiver.recv())
.await
.is_err()
);
}

#[rstest]
Expand All @@ -85,7 +89,11 @@ async fn non_validator_do_not_publish_ready(
})
.await;

assert!(p2p_receiver.try_recv().is_err(),);
assert!(
tokio::time::timeout(Duration::from_millis(10), p2p_receiver.recv())
.await
.is_err()
);
}

#[fixture]
Expand Down

0 comments on commit d695440

Please sign in to comment.