use std::sync::Arc;
use std::sync::Mutex;
use std::time::Duration;
use prost_types::FieldMask;
use proto::ledger_service_server::LedgerService;
use proto::ledger_service_server::LedgerServiceServer;
use proto::subscription_service_server::SubscriptionService;
use proto::subscription_service_server::SubscriptionServiceServer;
use proto::transaction_execution_service_server::TransactionExecutionService;
use proto::transaction_execution_service_server::TransactionExecutionServiceServer;
use sui_rpc::Client;
use sui_rpc::proto::sui::rpc::v2 as proto;
const CHECKPOINT_TICK: Duration = Duration::from_millis(100);
#[derive(Clone)]
struct MockServer {
digest: String,
digest_at_seq: u64,
exec_delay: Duration,
exec_error: Option<tonic::Status>,
shortcut_checkpoint: Option<u64>,
shortcut_from_lookup: usize,
lookups: Arc<Mutex<Vec<proto::GetTransactionRequest>>>,
}
impl MockServer {
fn new(digest: String) -> Self {
Self {
digest,
digest_at_seq: 10_000,
exec_delay: Duration::ZERO,
exec_error: None,
shortcut_checkpoint: None,
shortcut_from_lookup: 1,
lookups: Arc::default(),
}
}
}
#[tonic::async_trait]
impl TransactionExecutionService for MockServer {
async fn execute_transaction(
&self,
_request: tonic::Request<proto::ExecuteTransactionRequest>,
) -> Result<tonic::Response<proto::ExecuteTransactionResponse>, tonic::Status> {
tokio::time::sleep(self.exec_delay).await;
if let Some(status) = &self.exec_error {
return Err(status.clone());
}
Ok(tonic::Response::new(
proto::ExecuteTransactionResponse::default(),
))
}
}
#[tonic::async_trait]
impl LedgerService for MockServer {
async fn get_transaction(
&self,
request: tonic::Request<proto::GetTransactionRequest>,
) -> Result<tonic::Response<proto::GetTransactionResponse>, tonic::Status> {
let lookup_count = {
let mut lookups = self.lookups.lock().unwrap();
lookups.push(request.into_inner());
lookups.len()
};
let mut response = proto::GetTransactionResponse::default();
if let Some(checkpoint) = self.shortcut_checkpoint
&& lookup_count >= self.shortcut_from_lookup
{
let mut transaction = proto::ExecutedTransaction::default();
transaction.digest = Some(self.digest.clone());
transaction.checkpoint = Some(checkpoint);
response.transaction = Some(transaction);
}
Ok(tonic::Response::new(response))
}
}
#[tonic::async_trait]
impl SubscriptionService for MockServer {
async fn subscribe_checkpoints(
&self,
_request: tonic::Request<proto::SubscribeCheckpointsRequest>,
) -> Result<
tonic::Response<tonic::codegen::BoxStream<proto::SubscribeCheckpointsResponse>>,
tonic::Status,
> {
let digest = self.digest.clone();
let digest_at_seq = self.digest_at_seq;
let stream = futures::stream::unfold(0u64, move |seq| {
let digest = digest.clone();
async move {
tokio::time::sleep(CHECKPOINT_TICK).await;
let mut transaction = proto::ExecutedTransaction::default();
transaction.digest = Some(if seq == digest_at_seq {
digest
} else {
"filler".to_owned()
});
let mut checkpoint = proto::Checkpoint::default();
checkpoint.sequence_number = Some(seq);
checkpoint.transactions = vec![transaction];
let mut resp = proto::SubscribeCheckpointsResponse::default();
resp.cursor = Some(seq);
resp.checkpoint = Some(checkpoint);
Some((Ok(resp), seq + 1))
}
});
Ok(tonic::Response::new(Box::pin(stream)))
}
}
async fn spawn_mock_server(mock: MockServer) -> std::net::SocketAddr {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind mock server listener");
let addr = listener.local_addr().expect("mock server local addr");
tokio::spawn(async move {
tonic::transport::Server::builder()
.add_service(TransactionExecutionServiceServer::new(mock.clone()))
.add_service(LedgerServiceServer::new(mock.clone()))
.add_service(SubscriptionServiceServer::new(mock))
.serve_with_incoming(tonic::transport::server::TcpIncoming::from(listener))
.await
.expect("mock server exited with an error");
});
addr
}
fn test_transaction() -> (proto::Transaction, String) {
let transaction = sui_sdk_types::Transaction {
kind: sui_sdk_types::TransactionKind::ProgrammableTransaction(
sui_sdk_types::ProgrammableTransaction {
inputs: Vec::new(),
commands: Vec::new(),
},
),
sender: sui_sdk_types::Address::ZERO,
gas_payment: sui_sdk_types::GasPayment {
objects: Vec::new(),
owner: sui_sdk_types::Address::ZERO,
price: 1,
budget: 1,
},
expiration: sui_sdk_types::TransactionExpiration::None,
};
let digest = transaction.digest().to_string();
(proto::Transaction::from(transaction), digest)
}
#[tokio::test(flavor = "multi_thread")]
async fn subscription_survives_slow_execution() {
let (transaction, digest) = test_transaction();
let digest_at_seq = 50;
let addr = spawn_mock_server(MockServer {
digest_at_seq,
exec_delay: Duration::from_secs(4),
..MockServer::new(digest)
})
.await;
let mut client = Client::new(format!("http://{addr}"))
.expect("client")
.with_body_idle_timeout(Duration::from_millis(1500));
let request = proto::ExecuteTransactionRequest::default().with_transaction(transaction);
let response = client
.execute_transaction_and_wait_for_checkpoint(request, Duration::from_secs(30))
.await
.expect("execution and checkpoint confirmation");
assert_eq!(response.get_ref().transaction().checkpoint(), digest_at_seq);
}
#[tokio::test(flavor = "multi_thread")]
async fn already_checkpointed_transaction_short_circuits() {
let (transaction, digest) = test_transaction();
let addr = spawn_mock_server(MockServer {
shortcut_checkpoint: Some(7),
..MockServer::new(digest)
})
.await;
let mut client = Client::new(format!("http://{addr}")).expect("client");
let request = proto::ExecuteTransactionRequest::default().with_transaction(transaction);
let response = client
.execute_transaction_and_wait_for_checkpoint(request, Duration::from_secs(5))
.await
.expect("execution with an already checkpointed transaction");
assert_eq!(response.get_ref().transaction().checkpoint(), 7);
}
#[tokio::test(flavor = "multi_thread")]
async fn committed_duplicate_returns_before_execution_finishes() {
let (transaction, digest) = test_transaction();
let lookups = Arc::new(Mutex::new(Vec::new()));
let addr = spawn_mock_server(MockServer {
exec_delay: Duration::from_secs(60),
shortcut_checkpoint: Some(7),
lookups: lookups.clone(),
..MockServer::new(digest.clone())
})
.await;
let mut client = Client::new(format!("http://{addr}")).expect("client");
let request = proto::ExecuteTransactionRequest::default()
.with_transaction(transaction)
.with_read_mask(FieldMask {
paths: vec!["effects.status".to_owned()],
});
let start = std::time::Instant::now();
let response = client
.execute_transaction_and_wait_for_checkpoint(request, Duration::from_secs(5))
.await
.expect("committed duplicate resolved from the ledger probe");
assert!(
start.elapsed() < Duration::from_secs(30),
"the probe must resolve the call without waiting for execution"
);
assert_eq!(response.get_ref().transaction().checkpoint(), 7);
assert_eq!(response.get_ref().transaction().digest(), digest);
let probe_mask = lookups.lock().unwrap()[0]
.read_mask
.clone()
.expect("probe read mask");
for path in ["effects.status", "digest", "checkpoint", "timestamp"] {
assert!(
probe_mask.paths.iter().any(|p| p == path),
"probe read mask is missing `{path}`: {probe_mask:?}"
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn execution_error_with_committed_transaction_succeeds() {
let (transaction, digest) = test_transaction();
let addr = spawn_mock_server(MockServer {
exec_delay: Duration::from_millis(500),
exec_error: Some(tonic::Status::aborted("transaction already being executed")),
shortcut_checkpoint: Some(9),
shortcut_from_lookup: 2,
..MockServer::new(digest.clone())
})
.await;
let mut client = Client::new(format!("http://{addr}")).expect("client");
let request = proto::ExecuteTransactionRequest::default().with_transaction(transaction);
let response = client
.execute_transaction_and_wait_for_checkpoint(request, Duration::from_secs(5))
.await
.expect("execution error downgraded for a committed transaction");
assert_eq!(response.get_ref().transaction().checkpoint(), 9);
assert_eq!(response.get_ref().transaction().digest(), digest);
}
#[tokio::test(flavor = "multi_thread")]
async fn execution_error_without_commit_surfaces() {
let (transaction, digest) = test_transaction();
let addr = spawn_mock_server(MockServer {
exec_error: Some(tonic::Status::internal("boom")),
..MockServer::new(digest)
})
.await;
let mut client = Client::new(format!("http://{addr}")).expect("client");
let request = proto::ExecuteTransactionRequest::default().with_transaction(transaction);
let err = client
.execute_transaction_and_wait_for_checkpoint(request, Duration::from_secs(5))
.await
.expect_err("execution error with no committed transaction");
match err {
sui_rpc::client::ExecuteAndWaitError::RpcError(status) => {
assert_eq!(status.code(), tonic::Code::Internal);
}
other => panic!("expected RpcError, got {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn missing_transaction_fails_without_rpc() {
let mut client = Client::new("http://127.0.0.1:1").expect("client");
let err = client
.execute_transaction_and_wait_for_checkpoint(
proto::ExecuteTransactionRequest::default(),
Duration::from_secs(1),
)
.await
.expect_err("request without a transaction");
assert!(matches!(
err,
sui_rpc::client::ExecuteAndWaitError::MissingTransaction
));
}