Transfer Engine Rust API#
This page documents the Rust crate transfer_engine_rust (located at
mooncake-transfer-engine/rust). It is a library wrapper around the
Transfer Engine C API (transfer_engine_c.h).
Hot-path types (TransferRequest, BufferEntry, TransferStatus) are
#[repr(C)] and compile-time layout-checked against the bindgen C types, so
submit_transfer passes a Rust slice to C with no heap allocation and no
per-request copy on the Rust side.
For Transfer Engine design docs and non-Rust APIs, see:
Transfer Engine design docs:
design/transfer-engine/indexTransfer Engine C++ API:
api-reference/cpp/transfer-engineTransfer Engine Python API:
api-reference/python/transfer-engine
Build & runtime prerequisites#
Build:
Requires a Rust toolchain and libclang (bindgen).
CMake:
-DWITH_RUST_EXAMPLE=ON, thencmake --build build --target build_transfer_engine_rust.Or Cargo after exporting
MOONCAKE_BUILD_DIR/MOONCAKE_TE_LIB_DIR/MOONCAKE_TE_INCLUDE_DIR(seemooncake-transfer-engine/rust/README.md).
Runtime:
Dynamic linker must find Transfer Engine shared libraries (
libasio.so, …).A metadata server (HTTP metadata, etcd, or
P2PHANDSHAKE) must be reachable.GitHub Actions runs
scripts/ci/run_transfer_engine_rust_smoke.shafter the C++ build (cargo test --libplus the TCP loopbackminimal_smoketest).
Quick start#
use transfer_engine_rust::{MemoryPool, TransferEngine, TransferRequest, WILDCARD_LOCATION};
fn main() -> Result<(), transfer_engine_rust::EngineError> {
let engine = TransferEngine::initialize(
"127.0.0.1:12345",
"http://127.0.0.1:8080/metadata",
"tcp",
"",
)?;
let pool = MemoryPool::new(1 << 20);
unsafe {
engine.register_local_memory(pool.as_void_ptr(), pool.len(), WILDCARD_LOCATION)?;
let seg = engine.open_segment("peer:12345")?;
let req = TransferRequest::write(pool.as_void_ptr(), seg, /*offset*/ 0, 4096);
engine.submit_and_wait(&[req], None)?;
engine.unregister_local_memory(pool.as_void_ptr())?;
}
Ok(())
}
initialize(local_hostname, metadata_server, protocol, device_name) matches
the Python constructor. device_name is accepted for API compatibility; NIC
filtering is done with MC_TE_FILTERS because the C ABI has no device-name
argument.
The lower-level constructors map onto createTransferEngine:
TransferEngine::new(metadata_uri, local_server_name, rpc_port)TransferEngine::create(TransferEngineOptions { … })
Mental model#
Register local memory regions as RDMA/TCP-capable buffers.
Open a remote segment to obtain a
SegmentId.Allocate a
BatchIdfor a fixed number of requests.Submit a
&[TransferRequest](zero-copy FFI).Poll
get_transfer_status/wait_all, thenfree_batch_id.
Python-shaped helpers (transfer_sync_write, batch_transfer_sync_*,
transfer_submit_write) cache segment ids by hostname. They allocate a
batch internally. Use submit_transfer + wait_all when you need to keep
the batch/request arrays on the stack.
API reference#
Types#
Opcode::{Read, Write}—OPCODE_READ/OPCODE_WRITETransferStatusCode::{Waiting, Pending, Invalid, Canceled, Completed, Timeout, Failed}TransferRequest { opcode, source, target_id, target_offset, length }— layout matchestransfer_request_t. Helpers:TransferRequest::read,TransferRequest::write.BufferEntry { addr, length }— layout matchesbuffer_entry_tTransferStatus { status, transferred_bytes }— layout matchestransfer_status_tBatchId(u64)—INVALID_BATCHon allocate failureNotifyMsg { name, msg }NicLoadStat { device_name, inflight_bytes, ewma_bandwidth_bps }MemoryPool— page-aligned, zeroed host buffer for registrationWILDCARD_LOCATION("*"),LOCAL_SEGMENT(0)
Engine lifecycle#
initialize(local_hostname, metadata_server, protocol, device_name)new/creatediscover_topologyinstall_transport/uninstall_transportlocal_ip_and_portDrop destroys the native handle (no double-free)
Memory#
All pointer APIs are unsafe. Registered memory must stay valid until
unregistered.
register_local_memory/register_local_memory_ex(remote-accessible flag)unregister_local_memoryregister_memory/unregister_memory— Python aliases usingWILDCARD_LOCATIONregister_local_memory_batch(&[BufferEntry])— zero-copyunregister_local_memory_batch(&[*mut c_void])
Segments#
open_segment/open_segment_no_cache/open_segment_cachedclose_segmentwarmup_efa_segmentremove_local_segmentsync_segment_cache
Transfers (zero-copy hot path)#
allocate_batch_id(batch_size)submit_transfer(batch_id, &[TransferRequest])submit_transfer_with_notify(batch_id, requests, &NotifyMsg)get_transfer_status(batch_id, task_id) -> TransferStatuswait_all(batch_id, count, timeout)submit_and_wait(&[TransferRequest], timeout)— allocate + submit + wait + freefree_batch_id
Python-shaped transfers#
These open (and cache) a segment by hostname:
transfer_sync/transfer_sync_write/transfer_sync_readbatch_transfer_sync/batch_transfer_sync_write/batch_transfer_sync_readtransfer_submit_write— returnsBatchId; caller mustfree_batch_idtransfer_check_status— polls task 0; does not free the batch
Notifications and diagnostics#
take_notifies() -> Vec<NotifyMsg>send_notify(target_id, &NotifyMsg)nic_load_stats() -> Vec<NicLoadStat>enable_graceful_shutdownshow_links(json: bool) -> String
Errors#
EngineError (thiserror, #[non_exhaustive]):
NullHandleInvalidString(interior NUL)OperationFailed(i32)— raw C statusInvalidArgumentTransferFailedTimeout
Safety & thread-safety#
TransferEngineisSend + Sync; the C++ engine serializes internally.Pointer arguments must satisfy Rust aliasing and lifetime rules.
Registered memory must remain valid until unregistered.
submit_transferdoes not copy request bytes; do not mutate a submittedTransferRequestuntil the C call returns (the C layer copies into its own vector before returning).