Skip to main content

Crate msgtrans

Crate msgtrans 

Source
Expand description

ยง๐Ÿš€ MsgTrans - Modern Multi-Protocol Communication Framework

Rust License Crates.io Docs.rs

๐ŸŒ 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 Connection trait 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: Packet carries a Bytes payload, 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 Connection trait

ยง๐ŸŽฏ 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:

# 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
pluginDeprecated
protocol
stream
transport
Transport layer: client, server, session actors and request lifecycle.

Structsยง

SessionId
Type-safe wrapper for session ID

Type Aliasesยง

PacketId
Result