Expand description
ยง๐ MsgTrans - Modern Multi-Protocol Communication Framework
๐ Language: English | ็ฎไฝไธญๆ
Modern multi-protocol communication framework with a unified interface over TCP, WebSocket and QUIC
ยง๐ Core Features
ยง๐๏ธ Unified Architecture
- Three-layer architecture: Application โ Transport โ Protocol, with clear separation
- Protocol-agnostic business logic: one codebase, multi-protocol deployment
- Configuration-driven: switch protocols through configuration without changing business logic
- Pluggable adapters: implement the
Connectiontrait to add a new protocol
ยงโก Modern Concurrency
- Lock-free internals: per-session actors and lock-free maps avoid Mutex contention on the hot path
- Zero-copy packets:
Packetcarries aBytespayload, handed straight to the wire where possible - Event-driven model: fully asynchronous, non-blocking event handling
- Bounded backpressure: outbound queues are bounded per connection so a slow peer cannot exhaust memory or stall a fan-out loop
ยง๐ Protocols
- TCP - reliable stream transport
- WebSocket - real-time web communication
- QUIC - modern UDP-based transport
- Custom protocols - implement the
Connectiontrait
ยง๐ฏ Minimalist API
- Builder pattern: fluent, readable configuration
- Type safety: compile-time checked configuration
- Sensible defaults: works out of the box, tune only when needed
ยง๐ Quick Start
ยงInstallation
[dependencies]
msgtrans = "1.0"ยงCreate a Multi-Protocol Server
use msgtrans::{
transport::TransportServerBuilder,
protocol::{TcpServerConfig, WebSocketServerConfig, QuicServerConfig},
event::ServerEvent,
tokio,
};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// Configure multiple protocols - the same business logic serves all of them.
let tcp_config = TcpServerConfig::new("127.0.0.1:8001")?;
let websocket_config = WebSocketServerConfig::new("127.0.0.1:8002")?.with_path("/ws");
let quic_config = QuicServerConfig::new("127.0.0.1:8003")?;
let server = TransportServerBuilder::new()
.max_connections(10000)
.with_protocol(tcp_config)
.with_protocol(websocket_config)
.with_protocol(quic_config)
.build()
.await?;
// Subscribe to events before serving so nothing is missed.
let mut events = server.subscribe_events();
// Drive the listeners. `serve()` runs until the server is stopped, so spawn
// it and handle events on the main task.
let server_for_events = server.clone();
tokio::spawn(async move {
while let Ok(event) = events.recv().await {
match event {
ServerEvent::ConnectionEstablished { session_id, .. } => {
println!("New connection: {session_id}");
}
ServerEvent::MessageReceived { session_id, context } => {
// Echo the message back - protocol transparent.
let response = format!("Echo: {}", String::from_utf8_lossy(&context.data));
let _ = server_for_events.send(session_id, response.as_bytes()).await;
}
ServerEvent::ConnectionClosed { session_id, .. } => {
println!("Connection closed: {session_id}");
}
_ => {}
}
}
});
server.serve().await?;
Ok(())
}ยงCreate a Client Connection
use msgtrans::{
transport::TransportClientBuilder,
protocol::TcpClientConfig,
event::ClientEvent,
tokio,
};
use std::time::Duration;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let tcp_config = TcpClientConfig::new("127.0.0.1:8001")?
.with_connect_timeout(Duration::from_secs(30));
let mut client = TransportClientBuilder::new()
.with_protocol(tcp_config)
.build()
.await?;
client.connect().await?;
// Send a one-way message.
client.send("Hello, MsgTrans!".as_bytes()).await?;
println!("Message sent");
// Send a request and wait for the response.
let result = client.request("What time is it?".as_bytes()).await?;
if let Some(data) = result.data {
println!("Received response: {}", String::from_utf8_lossy(&data));
} else {
println!("Request timed out");
}
// Consume events.
let mut events = client.subscribe_events();
tokio::spawn(async move {
while let Ok(event) = events.recv().await {
match event {
ClientEvent::MessageReceived(context) => {
println!("Received: {}", String::from_utf8_lossy(&context.data));
}
ClientEvent::Disconnected { .. } => break,
_ => {}
}
}
});
Ok(())
}ยง๐๏ธ Architecture Design
ยงThree-Layer Architecture
+-------------------------------------+
| Application Layer | <- Business logic, protocol-agnostic
+-------------------------------------+
| Transport Layer | <- Connection management, unified API
| - TransportServer / TransportClient| - Connection lifecycle
| - SessionActor | - Event routing
| - RequestRegistry | - Request/response lifecycle
+-------------------------------------+
| Protocol Layer | <- Protocol implementation
| - TCP / WebSocket / QUIC adapters | - Connection trait
| - Protocol configs | - Protocol registration
+-------------------------------------+ยงDesign Principles
Unified abstraction, protocol transparency โ TransportServer/TransportClient
expose one business interface; each adapter implements the Connection trait and
hides protocol details.
Configuration-driven โ the same server code runs on any protocol; only the
config passed to .with_protocol(..) changes:
// TCP server
let server = TransportServerBuilder::new()
.with_protocol(TcpServerConfig::new("0.0.0.0:8080")?)
.build().await?;
// QUIC server - identical business logic
let server = TransportServerBuilder::new()
.with_protocol(QuicServerConfig::new("0.0.0.0:8080")?)
.build().await?;ยงEvent-Driven Model
use msgtrans::{
event::{ServerEvent, ClientEvent},
command::ConnectionInfo,
error::{TransportError, CloseReason},
SessionId, TransportContext,
};
// Server events
ServerEvent::ConnectionEstablished { session_id, info } => { /* ... */ }
ServerEvent::MessageReceived { session_id, context } => { /* ... */ }
ServerEvent::MessageSent { session_id, message_id } => { /* ... */ }
ServerEvent::ConnectionClosed { session_id, reason } => { /* ... */ }
ServerEvent::TransportError { session_id, error } => { /* ... */ }
// Client events
ClientEvent::Connected { info } => { /* ... */ }
ClientEvent::MessageReceived(context) => { /* ... */ }
ClientEvent::MessageSent { message_id } => { /* ... */ }
ClientEvent::Disconnected { reason } => { /* ... */ }
ClientEvent::Error { error } => { /* ... */ }ยงโก Usage Patterns
ยงConcurrent Sending
TransportServer is cheaply cloneable (it shares state via Arc), so it can be
moved into spawned tasks for concurrent, lock-free session access:
let mut events = server.subscribe_events();
while let Ok(event) = events.recv().await {
if let ServerEvent::MessageReceived { session_id, context } = event {
let server = server.clone();
tokio::spawn(async move {
let response = format!("Echo: {}", String::from_utf8_lossy(&context.data));
let _ = server.send(session_id, response.as_bytes()).await;
});
}
}ยงRequest / Response
let response = client.request(b"Get user data").await?;
if let Some(data) = response.data {
println!("Got {} bytes", data.len());
} else {
println!("Request timed out");
}ยง๐ Protocol Extension
To add a protocol, implement the Connection trait for your adapter and a
matching config type. See the built-in adapters::{tcp, websocket, quic} for
complete, working references; the outline below shows the shape:
use msgtrans::{connection::Connection, packet::Packet, error::TransportError};
pub struct MyAdapter { /* protocol-specific state */ }
#[async_trait::async_trait]
impl Connection for MyAdapter {
async fn send(&mut self, packet: Packet) -> Result<(), TransportError> { /* ... */ }
// ... remaining Connection methods
}ยง๐ Usage Examples
ยงWebSocket Server
use msgtrans::{
transport::TransportServerBuilder,
protocol::WebSocketServerConfig,
event::ServerEvent,
tokio,
};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let config = WebSocketServerConfig::new("127.0.0.1:8080")?.with_path("/chat");
let server = TransportServerBuilder::new()
.with_protocol(config)
.max_connections(1000)
.build()
.await?;
let mut events = server.subscribe_events();
let server_for_events = server.clone();
tokio::spawn(async move {
while let Ok(event) = events.recv().await {
if let ServerEvent::MessageReceived { session_id, context } = event {
let msg = String::from_utf8_lossy(&context.data);
let _ = server_for_events
.send(session_id, format!("You said: {msg}").as_bytes())
.await;
}
}
});
server.serve().await?;
Ok(())
}ยงQUIC Client
use msgtrans::{
transport::TransportClientBuilder,
protocol::QuicClientConfig,
tokio,
};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// Local / self-signed test server: skip certificate verification.
// In production, drop danger_skip_verification() and configure a real
// server name and CA instead.
let config = QuicClientConfig::new("127.0.0.1:8003")?
.with_server_name("localhost")
.danger_skip_verification();
let mut client = TransportClientBuilder::new()
.with_protocol(config)
.build()
.await?;
client.connect().await?;
for i in 0..1000u32 {
client.send(format!("message {i}").as_bytes()).await?;
}
println!("Done");
Ok(())
}ยง๐ ๏ธ Configuration Options
ยงServer Configuration
use msgtrans::protocol::{TcpServerConfig, WebSocketServerConfig, QuicServerConfig};
use std::time::Duration;
let tcp_config = TcpServerConfig::new("0.0.0.0:8001")?
.with_max_connections(10000)
.with_keepalive(Some(Duration::from_secs(60)))
.with_nodelay(true)
.with_reuse_addr(true);
let ws_config = WebSocketServerConfig::new("0.0.0.0:8002")?
.with_path("/api/ws")
.with_max_frame_size(1024 * 1024)
.with_max_connections(5000);
let quic_config = QuicServerConfig::new("0.0.0.0:8003")?
.with_cert_pem(std::fs::read_to_string("cert.pem")?)
.with_key_pem(std::fs::read_to_string("key.pem")?)
.with_max_concurrent_streams(1000);ยง๐ง Advanced Features
ยงStatistics
let active = server.session_count().await;
println!("Active sessions: {active}");ยงGraceful Error Handling
match client.send("Hello, World!".as_bytes()).await {
Ok(result) => println!("Sent (ID: {})", result.message_id),
Err(TransportError::Connection { .. }) => {
println!("Connection lost, reconnecting");
client.connect().await?;
}
Err(TransportError::Protocol { protocol, reason }) => {
println!("Protocol error [{protocol}]: {reason}");
}
Err(e) => println!("Other error: {e}"),
}ยงGraceful Shutdown
use msgtrans::{transport::TransportServerBuilder, protocol::TcpServerConfig};
use std::time::Duration;
let server = TransportServerBuilder::new()
.with_protocol(TcpServerConfig::new("0.0.0.0:8001")?)
.graceful_shutdown(Some(Duration::from_secs(30)))
.build().await?;
// ... later, drain active sessions using the configured timeout:
server.stop().await;ยง๐ Documentation and Examples
The examples/ directory contains complete, runnable programs:
echo_server.rs- multi-protocol echo serverecho_client_tcp.rs- TCP clientecho_client_websocket.rs- WebSocket clientecho_client_quic.rs- QUIC clientload_test.rs/load_test_server.rs- load testingpacket.rs- packet serialization
# Start the multi-protocol echo server
cargo run --example echo_server
# In another terminal, run a client
cargo run --example echo_client_tcpยง๐ Use Cases
- Game servers - high-concurrency real-time communication
- Chat systems - multi-protocol instant messaging
- Microservice communication - efficient inter-service transport
- Real-time data - financial, monitoring and telemetry systems
- IoT platforms - large-scale device connection management
- Protocol gateways - multi-protocol conversion and proxying
ยง๐ License
Licensed under the Apache License 2.0.
Copyright ยฉ 2024 zoujiaqing
ยง๐ค Contributing
Issues and Pull Requests are welcome.
Re-exportsยง
pub use command::ConnectionInfo;pub use command::TransportCommand;pub use command::TransportStats;pub use error::CloseReason;pub use error::TransportError;pub use event::ClientEvent;pub use event::QuicEvent;pub use event::TcpEvent;pub use event::TransportEvent;pub use event::WebSocketEvent;pub use packet::FramePolicy;pub use packet::Packet;pub use packet::PacketError;pub use packet::PacketType;pub use stream::ClientEventStream;pub use stream::EventStream;pub use stream::PacketStream;pub use transport::AcceptorConfig;pub use transport::BackpressureStrategy;pub use transport::CircuitBreakerConfig;pub use transport::ConnectionPool;Deprecated pub use transport::ConnectionPoolConfig;pub use transport::ExpertConfig;Deprecated pub use transport::LoadBalancerConfig;pub use transport::LockFreeCounter;pub use transport::LockFreeHashMap;pub use transport::LockFreeQueue;pub use transport::MemoryPool;pub use transport::MemoryStats;pub use transport::MemoryStatsSnapshot;pub use transport::PerformanceConfig;pub use transport::ProtocolStats;pub use transport::RateLimiterConfig;pub use transport::RetryConfig;pub use transport::SmartPoolConfig;pub use transport::Transport;pub use transport::TransportClient;pub use transport::TransportClientBuilder;pub use transport::TransportConfig;pub use transport::TransportContext;pub use transport::TransportServer;pub use transport::TransportServerBuilder;pub use protocol::ClientConfig;pub use protocol::QuicClientConfig;pub use protocol::QuicServerConfig;pub use protocol::ServerConfig;pub use protocol::TcpClientConfig;pub use protocol::TcpServerConfig;pub use protocol::WebSocketClientConfig;pub use protocol::WebSocketServerConfig;pub use connection::Connection;pub use connection::ConnectionFactory;pub use connection::Server;pub use plugin::PluginInfo;Deprecated pub use plugin::PluginManager;Deprecated pub use plugin::ProtocolPlugin;Deprecated pub use tokio;
Modulesยง
- adapters
- command
- connection
- error
- event
- packet
- plugin
Deprecated - protocol
- stream
- transport
- Transport layer: client, server, session actors and request lifecycle.
Structsยง
- Session
Id - Type-safe wrapper for session ID