feat(network): scaffold local state query server (#280)

This commit is contained in:
Harper 2023-10-26 19:02:11 +01:00 committed by GitHub
parent eace470b4f
commit fb1cc54b0c
9 changed files with 656 additions and 64 deletions

View file

@ -1,15 +1,23 @@
use std::fs;
use std::net::{Ipv4Addr, SocketAddrV4};
use std::time::Duration;
use pallas_network::facades::{PeerClient, PeerServer};
use pallas_network::facades::{NodeClient, PeerClient, PeerServer};
use pallas_network::miniprotocols::blockfetch::BlockRequest;
use pallas_network::miniprotocols::handshake::n2c;
use pallas_network::miniprotocols::handshake::n2n::VersionData;
use pallas_network::miniprotocols::localstate::queries::{GenericResponse, Request};
use pallas_network::miniprotocols::localstate::{ClientAcquireRequest, ClientQueryRequest};
use pallas_network::miniprotocols::chainsync::{ClientRequest, HeaderContent, Tip};
use pallas_network::miniprotocols::{
blockfetch,
chainsync::{self, NextResponse},
Point,
};
use tokio::net::TcpListener;
use pallas_network::miniprotocols::{handshake, localstate};
use pallas_network::multiplexer::{Bearer, Plexer};
use std::path::Path;
use tokio::net::{TcpListener, UnixListener};
#[tokio::test]
#[ignore]
@ -247,6 +255,150 @@ pub async fn blockfetch_server_and_client_happy_path() {
_ = tokio::join!(client, server);
}
#[tokio::test]
#[ignore]
pub async fn local_state_query_server_and_client_happy_path() {
let server = tokio::spawn({
async move {
// server setup
let socket_path = Path::new("node.socket");
if socket_path.exists() {
fs::remove_file(&socket_path).unwrap();
}
let unix_listener = UnixListener::bind(socket_path).unwrap();
let (bearer, _) = Bearer::accept_unix(&unix_listener).await.unwrap();
let mut server_plexer = Plexer::new(bearer);
let mut server_hs: handshake::Server<n2c::VersionData> =
handshake::Server::new(server_plexer.subscribe_server(0));
let mut server_sq: localstate::Server =
localstate::Server::new(server_plexer.subscribe_server(7));
tokio::spawn(async move { server_plexer.run().await });
server_hs.receive_proposed_versions().await.unwrap();
server_hs
.accept_version(10, n2c::VersionData::new(0, Some(false)))
.await
.unwrap();
// server receives range from client, sends blocks
let ClientAcquireRequest(maybe_point) =
server_sq.recv_while_idle().await.unwrap().unwrap();
assert_eq!(maybe_point, Some(Point::Origin));
assert_eq!(*server_sq.state(), localstate::State::Acquiring);
// server_bf.send_block_range(bodies).await.unwrap();
server_sq.send_acquired().await.unwrap();
assert_eq!(*server_sq.state(), localstate::State::Acquired);
// server receives query from client
let query = match server_sq.recv_while_acquired().await.unwrap() {
ClientQueryRequest::Query(q) => q,
x => panic!("unexpected message from client: {x:?}"),
};
assert_eq!(
query,
Request::BlockQuery(localstate::queries::BlockQuery::GetStakePools)
);
assert_eq!(*server_sq.state(), localstate::State::Querying);
server_sq
.send_result(GenericResponse::new(hex::decode("82011A008BD423").unwrap()))
.await
.unwrap();
assert_eq!(*server_sq.state(), localstate::State::Acquired);
// server receives reaquire from the client
let maybe_point = match server_sq.recv_while_acquired().await.unwrap() {
ClientQueryRequest::ReAcquire(p) => p,
x => panic!("unexpected message from client: {x:?}"),
};
assert_eq!(maybe_point, Some(Point::Specific(1337, vec![1, 2, 3])));
assert_eq!(*server_sq.state(), localstate::State::Acquiring);
server_sq.send_acquired().await.unwrap();
// server receives release from the client
match server_sq.recv_while_acquired().await.unwrap() {
ClientQueryRequest::Release => (),
x => panic!("unexpected message from client: {x:?}"),
};
assert!(server_sq.recv_while_idle().await.unwrap().is_none());
assert_eq!(*server_sq.state(), localstate::State::Done);
}
});
let client = tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(1)).await;
// client setup
let socket_path = "node.socket";
let mut client_to_server_conn = NodeClient::connect(socket_path, 0).await.unwrap();
let client_sq = client_to_server_conn.statequery();
// client sends acquire
client_sq.send_acquire(Some(Point::Origin)).await.unwrap();
client_sq.recv_while_acquiring().await.unwrap();
assert_eq!(*client_sq.state(), localstate::State::Acquired);
// client sends a BlockQuery
client_sq
.send_query(Request::BlockQuery(
localstate::queries::BlockQuery::GetStakePools,
))
.await
.unwrap();
let resp = client_sq.recv_while_querying().await.unwrap();
assert_eq!(
resp,
GenericResponse::new(hex::decode("82011A008BD423").unwrap())
);
// client sends a ReAquire
client_sq
.send_reacquire(Some(Point::Specific(1337, vec![1, 2, 3])))
.await
.unwrap();
client_sq.recv_while_acquiring().await.unwrap();
client_sq.send_release().await.unwrap();
client_sq.send_done().await.unwrap();
});
_ = tokio::join!(client, server);
}
#[tokio::test]
#[ignore]
pub async fn chainsync_server_and_client_happy_path_n2n() {
@ -263,9 +415,21 @@ pub async fn chainsync_server_and_client_happy_path_n2n() {
.await
.unwrap();
let mut peer_server = PeerServer::accept(&server_listener, 0).await.unwrap();
let (bearer, _) = Bearer::accept_tcp(&server_listener).await.unwrap();
let server_cs = peer_server.chainsync();
let mut server_plexer = Plexer::new(bearer);
let mut server_hs: handshake::Server<VersionData> =
handshake::Server::new(server_plexer.subscribe_server(0));
let mut server_cs = chainsync::N2NServer::new(server_plexer.subscribe_server(2));
tokio::spawn(async move { server_plexer.run().await });
server_hs.receive_proposed_versions().await.unwrap();
server_hs
.accept_version(10, VersionData::new(0, false))
.await
.unwrap();
// server receives find intersect from client, sends intersect point
@ -359,7 +523,7 @@ pub async fn chainsync_server_and_client_happy_path_n2n() {
});
let client = tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(1)).await;
tokio::time::sleep(Duration::from_secs(2)).await;
// client setup
@ -432,6 +596,4 @@ pub async fn chainsync_server_and_client_happy_path_n2n() {
});
_ = tokio::join!(client, server);
}
// TODO: redo txsubmission client test
}