Skip to main content

moirai_transport/
lib.rs

1//! Unified transport layer for Moirai concurrency library.
2//!
3//! This module provides transport abstractions that work across different
4//! communication boundaries: threads, processes, and machines. It builds on
5//! top of the core channel primitives to provide location-transparent messaging.
6//!
7//! # Design Principles
8//! - Location transparency: same API for local and remote communication
9//! - Zero-copy optimization for local transport
10//! - Pluggable transport backends (in-memory, IPC, network)
11//! - Integration with Moirai scheduler for optimal performance
12
13#![deny(missing_docs)]
14
15#[cfg(any(unix, windows))]
16mod ipc;
17mod network;
18mod router;
19mod transport;
20
21pub mod payload;
22pub mod process;
23pub mod remote_task;
24#[cfg(feature = "scheduler-routes")]
25pub mod route;
26pub mod safe_channel;
27
28use std::sync::{Mutex, MutexGuard, PoisonError};
29
30/// Crate-wide lock policy: recover from poisoning instead of propagating the
31/// panic. Guarded state here (channel maps, subscription lists, connection
32/// states) stays structurally valid under a poisoned lock — a writer that
33/// panicked mid-critical-section cannot leave a torn invariant in these maps —
34/// so continuing with the recovered guard is sound. Matches the pal reactor
35/// backends' `lock_mutex` helpers.
36pub(crate) fn lock_mutex<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
37    mutex.lock().unwrap_or_else(PoisonError::into_inner)
38}
39
40// Re-export core channel types for compatibility
41/// Shared-memory same-machine IPC transport (Unix/Windows only).
42#[cfg(any(unix, windows))]
43pub use ipc::IpcTransport;
44pub use moirai_core::channel::{
45    ChannelError as TransportError, MpmcReceiver as Receiver, MpmcSender as Sender,
46};
47#[cfg(feature = "network")]
48pub use network::TcpTransport;
49pub(crate) use network::{NETWORK_IO_TIMEOUT, read_network_frame_from_stream};
50pub use network::{NetworkListener, NetworkTransport};
51pub use router::{MessageRouter, RemoteAddress};
52// The canonical typed cross-boundary channel: rkyv-style archive serialization
53// over a transport (zero-copy borrowed views on receive).
54pub use safe_channel::{
55    ArchiveSerialize, ArchiveView, ArchivedMessage, ArchivedUniversalReceiver,
56    ArchivedUniversalSender,
57};
58pub use transport::{
59    Address, ConnectionManager, ConnectionState, InMemoryTransport, TransportManager,
60};
61
62/// Result type for transport operations
63pub type TransportResult<T> = Result<T, TransportError>;
64
65/// Transport trait for different communication mechanisms
66pub trait Transport: Send + Sync {
67    /// Send a message to the specified address
68    fn send(&self, target: &Address, data: Vec<u8>) -> TransportResult<()>;
69
70    /// Receive a message from the specified address
71    fn recv(&self, source: &Address) -> TransportResult<Vec<u8>>;
72
73    /// Check if the transport supports the given address
74    fn supports(&self, address: &Address) -> bool;
75}
76
77// A typed cross-boundary channel over a transport is provided by the rkyv-style
78// archive channels in `safe_channel` (`ArchivedUniversalSender<T: ArchiveSerialize>`
79// / `ArchivedUniversalReceiver<T: ArchiveView>`), re-exported below. The previous
80// `UniversalChannel<T: Send>` / `UniversalSender` / `UniversalReceiver` were
81// non-functional placeholders (their `send`/`recv` ignored their argument and
82// returned `Closed`): a channel generic over an arbitrary `Send` `T` cannot
83// serialize the value for transport without a serialization bound, which is
84// exactly what the archive traits add. They were removed in favor of the working
85// archive channels rather than left as mocks.