Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion example/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ ctrlc = { version = "3.0", features = ["termination"] }
tokio = { version = "1.0.1", features = ["signal", "time"] }
async-trait = "0.1.42"
rand = "0.8.5"

clap = { version = "4.5.40", features = ["derive"] }

[[example]]
name = "client"
Expand Down
4 changes: 3 additions & 1 deletion example/async-client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,9 @@ use ttrpc::r#async::Client;

#[tokio::main(flavor = "current_thread")]
async fn main() {
let c = Client::connect(utils::SOCK_ADDR).await.unwrap();
let sock_addr = utils::get_sock_addr();
let c = Client::connect(sock_addr).await.unwrap();

let hc = health_ttrpc::HealthClient::new(c.clone());
let ac = agent_ttrpc::AgentServiceClient::new(c);

Expand Down
5 changes: 3 additions & 2 deletions example/async-server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,10 +58,11 @@ async fn main() {
let hservice = health_ttrpc::create_health(Arc::new(HealthService {}));
let aservice = agent_ttrpc::create_agent_service(Arc::new(AgentService {}));

utils::remove_if_sock_exist(utils::SOCK_ADDR).unwrap();
let sock_addr = utils::get_sock_addr();
utils::remove_if_sock_exist(sock_addr).unwrap();

let mut server = Server::new()
.bind(utils::SOCK_ADDR)
.bind(sock_addr)
.unwrap()
.register_service(hservice)
.register_service(aservice);
Expand Down
4 changes: 3 additions & 1 deletion example/async-stream-client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,9 @@ use ttrpc::r#async::Client;
async fn main() {
simple_logging::log_to_stderr(log::LevelFilter::Info);

let c = Client::connect(utils::SOCK_ADDR).await.unwrap();
let sock_addr = utils::get_sock_addr();
let c = Client::connect(sock_addr).await.unwrap();

let sc = streaming_ttrpc::StreamingClient::new(c);

let _now = std::time::Instant::now();
Expand Down
6 changes: 4 additions & 2 deletions example/async-stream-server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -171,10 +171,12 @@ impl streaming_ttrpc::Streaming for StreamingService {
async fn main() {
simple_logging::log_to_stderr(LevelFilter::Info);
let service = streaming_ttrpc::create_streaming(Arc::new(StreamingService {}));
utils::remove_if_sock_exist(utils::SOCK_ADDR).unwrap();

let sock_addr = utils::get_sock_addr();
utils::remove_if_sock_exist(sock_addr).unwrap();

let mut server = Server::new()
.bind(utils::SOCK_ADDR)
.bind(sock_addr)
.unwrap()
.register_service(service);

Expand Down
4 changes: 3 additions & 1 deletion example/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,9 @@ fn main() {
fn connect_once() {
simple_logging::log_to_stderr(LevelFilter::Trace);

let c = Client::connect(utils::SOCK_ADDR).unwrap();
let sock_addr = utils::get_sock_addr();
let c = Client::connect(sock_addr).unwrap();

let hc = health_ttrpc::HealthClient::new(c.clone());
let ac = agent_ttrpc::AgentServiceClient::new(c);

Expand Down
6 changes: 4 additions & 2 deletions example/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,9 +58,11 @@ fn main() {
let hservice = health_ttrpc::create_health(Arc::new(HealthService {}));
let aservice = agent_ttrpc::create_agent_service(Arc::new(AgentService {}));

utils::remove_if_sock_exist(utils::SOCK_ADDR).unwrap();
let sock_addr = utils::get_sock_addr();
utils::remove_if_sock_exist(sock_addr).unwrap();

let mut server = ttrpc::Server::new()
.bind(utils::SOCK_ADDR)
.bind(sock_addr)
.unwrap()
.register_service(hservice)
.register_service(aservice);
Expand Down
32 changes: 30 additions & 2 deletions example/utils.rs
Original file line number Diff line number Diff line change
@@ -1,14 +1,42 @@
#![allow(dead_code)]
use clap::Parser;
use log::warn;
use std::io::Result;

#[cfg(unix)]
pub const SOCK_ADDR: &str = r"unix:///tmp/ttrpc-test";
pub const SOCK_ADDR_LOCAL: &str = r"unix:///tmp/ttrpc-test";

#[cfg(windows)]
pub const SOCK_ADDR: &str = r"\\.\pipe\ttrpc-test";
pub const SOCK_ADDR_LOCAL: &str = r"\\.\pipe\ttrpc-test";

pub const SOCK_ADDR_TCP: &str = r"tcp://127.0.0.1:65500";

#[derive(Debug, Default, Parser)]
pub struct Cli {
#[arg(long = "tcp")]
#[arg(help = "Use a TCP socket instead of a local one")]
pub tcp: bool,
}

pub fn get_sock_addr() -> &'static str {
let cli = Cli::parse();
if cli.tcp {
if cfg!(windows) {
warn!("'--tcp' flag ignored; TCP sockets not supported on Windows");
return SOCK_ADDR_LOCAL;
}
SOCK_ADDR_TCP
} else {
SOCK_ADDR_LOCAL
}
}

#[cfg(unix)]
pub fn remove_if_sock_exist(sock_addr: &str) -> Result<()> {
if sock_addr.starts_with("tcp://") {
return Ok(());
}

let path = sock_addr
.strip_prefix("unix://")
.expect("socket address is not expected");
Expand Down
9 changes: 9 additions & 0 deletions src/asynchronous/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,15 @@ impl Server {
Ok(self.add_listener(listener))
}

#[cfg(unix)]
/// # Safety
/// The file descriptor must represent a unix listener.
pub unsafe fn add_tcp_listener(self, fd: RawFd) -> Result<Server> {
let listener = Listener::from_raw_tcp_listener_fd(fd)
.map_err(err_to_others_err!(e, "from_raw_tcp_listener_fd error"))?;
Ok(self.add_listener(listener))
}

#[cfg(any(target_os = "linux", target_os = "android"))]
/// # Safety
/// The file descriptor must represent a vsock listener.
Expand Down
13 changes: 13 additions & 0 deletions src/asynchronous/transport/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,9 @@ macro_rules! io_other {
#[cfg(unix)]
mod unix;

#[cfg(unix)]
mod tcp;

#[cfg(any(target_os = "linux", target_os = "android"))]
mod vsock;

Expand All @@ -43,6 +46,11 @@ impl Listener {
return Self::bind_unix(addr);
}

#[cfg(unix)]
if let Some(addr) = addr.strip_prefix("tcp://") {
return Self::bind_tcp(addr);
}

#[cfg(any(target_os = "linux", target_os = "android"))]
if let Some(addr) = addr.strip_prefix("vsock://") {
return Self::bind_vsock(addr);
Expand Down Expand Up @@ -70,6 +78,11 @@ impl Socket {
return Self::connect_unix(addr).await;
}

#[cfg(unix)]
if let Some(addr) = addr.strip_prefix("tcp://") {
return Self::connect_tcp(addr).await;
}

#[cfg(any(target_os = "linux", target_os = "android"))]
if let Some(addr) = addr.strip_prefix("vsock://") {
return Self::connect_vsock(addr).await;
Expand Down
80 changes: 80 additions & 0 deletions src/asynchronous/transport/tcp.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
use std::convert::TryFrom;
use std::io::{Error as IoError, Result as IoResult};
use std::os::fd::{FromRawFd as _, RawFd};
use std::net::{
SocketAddr, TcpListener as StdTcpListener, TcpStream as StdTcpStream,
};

use async_stream::stream;
use tokio::net::{TcpListener, TcpStream};

use super::{Listener, Socket};

impl Listener {
pub fn bind_tcp(addr: impl AsRef<str>) -> IoResult<Self> {
let addr = parse_tcp_addr(addr)?;
let listener = StdTcpListener::bind(addr)?;
Self::try_from(listener)
}

/// # Safety
/// The file descriptor must represent a tcp listener.
pub unsafe fn from_raw_tcp_listener_fd(fd: std::os::fd::RawFd) -> IoResult<Self> {
let listener = unsafe { StdTcpListener::from_raw_fd(fd) };
Self::try_from(listener)
}
}

impl Socket {
pub async fn connect_tcp(addr: impl AsRef<str>) -> IoResult<Self> {
let addr = parse_tcp_addr(addr)?;
let socket = StdTcpStream::connect(addr)?;
Self::try_from(socket)
}

/// # Safety
/// The file descriptor must represent a tcp socket.
pub unsafe fn from_raw_tcp_socket_fd(fd: RawFd) -> IoResult<Self> {
let socket = unsafe { StdTcpStream::from_raw_fd(fd) };
Self::try_from(socket)
}
}

impl From<TcpListener> for Listener {
fn from(listener: TcpListener) -> Self {
Self::new(stream! {
loop {
yield listener.accept().await.map(|(socket, _)| socket);
}
})
}
}

impl TryFrom<StdTcpListener> for Listener {
type Error = IoError;
fn try_from(listener: StdTcpListener) -> IoResult<Self> {
listener.set_nonblocking(true)?;
Ok(Self::from(TcpListener::from_std(listener)?))
}
}

impl From<TcpStream> for Socket {
fn from(socket: TcpStream) -> Self {
Self::new(socket)
}
}

impl TryFrom<StdTcpStream> for Socket {
type Error = IoError;
fn try_from(socket: StdTcpStream) -> IoResult<Self> {
socket.set_nonblocking(true)?;
Ok(Self::from(TcpStream::from_std(socket)?))
}
}

fn parse_tcp_addr(addr: impl AsRef<str>) -> IoResult<SocketAddr> {
let addr = addr.as_ref();

addr.parse::<SocketAddr>()
.map_err(|e| io_other!("Failed to parse TCP address '{}': {}", addr, e))
}
Loading