use std::sync::Arc;
use anyhow::Result;
pub mod discovery;
pub mod events;
pub mod messenger;
pub mod observability;
pub mod queue;
pub mod rendezvous;
pub mod streaming;
pub mod transports;
#[cfg(feature = "simulation")]
pub mod simulation;
pub use velo_ext::{
AdmissionState, InstanceId, PeerInfo, ShutdownPolicy, Transport, WorkerAddress, WorkerId,
};
pub use velo_ext as ext;
pub use crate::messenger::{
Admitted, AmHandlerBuilder, AmSendBuilder, AmSyncBuilder, AsyncExecutor, Context, DispatchMode,
FireResult, Handler, HandlerExecutor, Messenger, MessengerBuilder, OrderedConfig, OrderingKey,
OverflowPolicy, PeerDiscovery, SyncExecutor, SyncResult, TypedContext, TypedUnaryBuilder,
TypedUnaryHandlerBuilder, TypedUnaryResult, UnaryBuilder, UnaryHandlerBuilder, UnaryResult,
UnifiedResponse, VeloEvents,
};
pub use crate::events::{
Event, EventAwaiter, EventBackend, EventHandle, EventManager, EventPoison, EventStatus,
};
pub use crate::streaming::{
AnchorManager, AttachError, SendError, StreamAnchor, StreamAnchorHandle, StreamController,
StreamError, StreamFrame, StreamSender,
};
pub use crate::rendezvous::{
DataHandle, DataMetadata, RegisterOptions, RendezvousManager, RendezvousWrite, StageMode,
};
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub use crate::rendezvous::rdma::{
Deregistered, PinnedBuf, RdmaConfig, RdmaError, RdmaPoolConfig, RdmaRendezvousConfig,
RegionGuard, RegionWatch, RegisterOwnedError,
};
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub use crate::rendezvous::write::PinnedWriter;
pub use crate::rendezvous::write::RdmaDestination;
pub use crate::observability::VeloMetrics;
#[derive(Debug, Clone)]
pub struct TcpConfig {
pub bind_addr: std::net::IpAddr,
}
impl Default for TcpConfig {
fn default() -> Self {
Self {
bind_addr: std::net::IpAddr::V4(std::net::Ipv4Addr::UNSPECIFIED),
}
}
}
impl TcpConfig {
pub fn new(bind_addr: std::net::IpAddr) -> Self {
Self { bind_addr }
}
}
#[cfg(feature = "grpc")]
#[derive(Debug, Clone)]
pub struct GrpcConfig {
pub bind_addr: std::net::SocketAddr,
}
#[cfg(feature = "grpc")]
impl Default for GrpcConfig {
fn default() -> Self {
Self {
bind_addr: "0.0.0.0:0".parse().unwrap(),
}
}
}
#[derive(Debug, Clone)]
pub enum StreamConfig {
Tcp(Option<TcpConfig>),
#[cfg(feature = "grpc")]
Grpc(Option<GrpcConfig>),
}
#[derive(Clone)]
pub struct Velo {
messenger: Arc<Messenger>,
anchor_manager: Arc<crate::streaming::AnchorManager>,
rendezvous_manager: Arc<crate::rendezvous::RendezvousManager>,
stream_transport: Arc<dyn crate::streaming::FrameTransport>,
#[cfg(all(target_os = "linux", feature = "ucx"))]
rdma: Option<Arc<crate::rendezvous::rdma::RdmaRegistry>>,
shutdown: Arc<ShutdownOnce>,
}
struct ShutdownOnce {
lock: tokio::sync::Mutex<()>,
done: std::sync::atomic::AtomicBool,
}
pub struct VeloBuilder {
inner: MessengerBuilder,
stream_config: Option<StreamConfig>,
mux_config: Option<crate::streaming::MuxConfig>,
metrics: Option<Arc<VeloMetrics>>,
#[cfg(all(target_os = "linux", feature = "ucx"))]
ucx_transport: Option<Arc<crate::transports::ucx::UcxTransport>>,
#[cfg(all(target_os = "linux", feature = "ucx"))]
rdma_config: Option<crate::rendezvous::rdma::RdmaConfig>,
}
impl VeloBuilder {
pub fn new() -> Self {
Self {
inner: MessengerBuilder::new(),
stream_config: None,
mux_config: None,
metrics: None,
#[cfg(all(target_os = "linux", feature = "ucx"))]
ucx_transport: None,
#[cfg(all(target_os = "linux", feature = "ucx"))]
rdma_config: None,
}
}
pub fn add_transport(mut self, transport: Arc<dyn Transport>) -> Self {
self.inner = self.inner.add_transport(transport);
self
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub fn add_ucx_transport(
mut self,
transport: Arc<crate::transports::ucx::UcxTransport>,
) -> Self {
self.ucx_transport = Some(Arc::clone(&transport));
self.add_transport(transport)
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub fn rdma_config(mut self, config: crate::rendezvous::rdma::RdmaConfig) -> Self {
self.rdma_config = Some(config);
self
}
pub fn stream_config(mut self, config: StreamConfig) -> Result<Self> {
if self.stream_config.is_some() {
return Err(anyhow::anyhow!(
"stream_config called more than once: only one streaming server allowed per Velo instance"
));
}
self.stream_config = Some(config);
Ok(self)
}
pub fn stream_bind_addr(self, addr: std::net::IpAddr) -> Self {
self.stream_config(StreamConfig::Tcp(Some(TcpConfig::new(addr))))
.unwrap()
}
pub fn messenger_mux(mut self, config: crate::streaming::MuxConfig) -> Result<Self> {
if self.mux_config.is_some() {
return Err(anyhow::anyhow!(
"messenger_mux called more than once: only one messenger mux is allowed per Velo instance"
));
}
self.mux_config = Some(config);
Ok(self)
}
pub fn discovery(mut self, discovery: Arc<dyn PeerDiscovery>) -> Self {
self.inner = self.inner.discovery(discovery);
self
}
pub fn metrics(mut self, metrics: Arc<VeloMetrics>) -> Self {
self.inner = self.inner.metrics(metrics.clone());
self.metrics = Some(metrics);
self
}
pub async fn build(self) -> Result<Arc<Velo>> {
let messenger = self.inner.build().await?;
let worker_id = messenger.instance_id().worker_id();
let resolved = self.stream_config.unwrap_or(StreamConfig::Tcp(None));
let stream_transport: Arc<dyn crate::streaming::FrameTransport> = match resolved {
StreamConfig::Tcp(tcp_cfg) => {
let bind_addr = tcp_cfg
.map(|c| c.bind_addr)
.unwrap_or(std::net::IpAddr::V4(std::net::Ipv4Addr::UNSPECIFIED));
let tcp = crate::streaming::TcpFrameTransport::new(bind_addr).await?;
if let Some(m) = self.metrics.as_ref() {
tcp.set_metrics(Arc::clone(m));
}
tcp as _
}
#[cfg(feature = "grpc")]
StreamConfig::Grpc(grpc_cfg) => {
let bind_addr = grpc_cfg
.map(|c| c.bind_addr)
.unwrap_or_else(|| "0.0.0.0:0".parse().unwrap());
let grpc = crate::streaming::GrpcFrameTransport::new(bind_addr)
.await
.map_err(|e| {
anyhow::anyhow!("Failed to start gRPC streaming transport: {}", e)
})?;
if let Some(m) = self.metrics.as_ref() {
grpc.set_metrics(Arc::clone(m));
}
grpc as _
}
};
let mut registry: std::collections::HashMap<
String,
Arc<dyn crate::streaming::FrameTransport>,
> = std::collections::HashMap::new();
registry.insert(
stream_transport.key().as_str().to_string(),
Arc::clone(&stream_transport),
);
let mux = match self.mux_config.filter(|config| config.enabled) {
Some(config) => {
let mux = crate::streaming::messenger_mux::MessengerMuxTransport::new(
Arc::clone(&messenger),
config,
self.metrics.clone(),
)?;
let mux_key = crate::streaming::FrameTransport::key(mux.as_ref());
registry.insert(
mux_key.as_str().to_string(),
Arc::clone(&mux) as Arc<dyn crate::streaming::FrameTransport>,
);
Some(mux)
}
None => None,
};
let anchor_manager = Arc::new(
crate::streaming::AnchorManagerBuilder::default()
.worker_id(worker_id)
.transport(Arc::clone(&stream_transport))
.transport_registry(Arc::new(registry))
.messenger(Some(Arc::clone(&messenger)))
.metrics(self.metrics.clone())
.build()
.map_err(|e| anyhow::anyhow!("{}", e))?,
);
if let Some(mux) = mux {
anchor_manager.install_mux(mux)?;
}
anchor_manager.register_handlers(Arc::clone(&messenger))?;
let rendezvous_manager = Arc::new(match self.metrics.as_ref() {
Some(m) => crate::rendezvous::RendezvousManager::with_metrics(worker_id, Arc::clone(m)),
None => crate::rendezvous::RendezvousManager::new(worker_id),
});
rendezvous_manager.register_handlers(Arc::clone(&messenger))?;
let stager = Arc::new(crate::rendezvous::RendezvousStager::new(Arc::clone(
&rendezvous_manager,
)));
let resolver = Arc::new(crate::rendezvous::RendezvousResolver::new(Arc::clone(
&rendezvous_manager,
)));
messenger.set_large_payload_support(stager, resolver);
#[cfg(all(target_os = "linux", feature = "ucx"))]
let rdma = match self.ucx_transport.as_ref() {
Some(transport) => {
let mut config = self.rdma_config.clone().unwrap_or_default();
if rdma_rendezvous_disabled_by_env() {
tracing::info!(
"VELO_RDMA_RENDEZVOUS_DISABLE is set: the rendezvous RDMA path is off. \
Staged data is still readable — every slot answers the chunked path."
);
config.rendezvous.enabled = false;
}
let rendezvous_config = config.rendezvous.clone();
let registry = Arc::new(crate::rendezvous::rdma::RdmaRegistry::new(
crate::rendezvous::rdma::UcxBackend::new(transport.rdma_endpoint()),
config,
messenger.runtime().clone(),
self.metrics.clone(),
));
rendezvous_manager.set_rdma_context(
Arc::clone(®istry),
rendezvous_config,
messenger.runtime(),
)?;
Some(registry)
}
None => None,
};
Ok(Arc::new(Velo {
messenger,
anchor_manager,
rendezvous_manager,
stream_transport,
#[cfg(all(target_os = "linux", feature = "ucx"))]
rdma,
shutdown: Arc::new(ShutdownOnce {
lock: tokio::sync::Mutex::new(()),
done: std::sync::atomic::AtomicBool::new(false),
}),
}))
}
}
impl Default for VeloBuilder {
fn default() -> Self {
Self::new()
}
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
fn rdma_rendezvous_disabled_by_env() -> bool {
rdma_rendezvous_disabled(
std::env::var("VELO_RDMA_RENDEZVOUS_DISABLE")
.ok()
.as_deref(),
)
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
fn rdma_rendezvous_disabled(value: Option<&str>) -> bool {
value.is_some_and(|v| {
let v = v.trim().to_ascii_lowercase();
v == "1" || v == "true" || v == "yes" || v == "on"
})
}
impl Velo {
pub fn builder() -> VeloBuilder {
VeloBuilder::new()
}
pub fn messenger(&self) -> &Arc<Messenger> {
&self.messenger
}
pub fn begin_drain(&self) {
self.messenger.begin_drain();
}
pub async fn graceful_shutdown(&self, policy: ShutdownPolicy) {
let _running = self.shutdown.lock.lock().await;
if self
.shutdown
.done
.load(std::sync::atomic::Ordering::Acquire)
{
return;
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
let policy = {
let started = std::time::Instant::now();
if let Some(rdma) = &self.rdma {
self.begin_drain();
self.rendezvous_manager.shutdown();
let budget = match &policy {
ShutdownPolicy::Timeout(deadline) => *deadline,
ShutdownPolicy::WaitForever => rdma.shutdown_timeout(),
};
rdma.shutdown(budget).await;
}
match &policy {
ShutdownPolicy::Timeout(deadline) => {
ShutdownPolicy::Timeout(deadline.saturating_sub(started.elapsed()))
}
ShutdownPolicy::WaitForever => ShutdownPolicy::WaitForever,
}
};
self.messenger.graceful_shutdown(policy).await;
#[cfg(all(target_os = "linux", feature = "ucx"))]
if let Some(rdma) = &self.rdma {
rdma.latch_all_deregistered();
}
self.shutdown
.done
.store(true, std::sync::atomic::Ordering::Release);
}
pub fn instance_id(&self) -> InstanceId {
self.messenger.instance_id()
}
pub fn flush_batch(&self) {
self.anchor_manager.flush_mux_batches();
}
pub fn peer_info(&self) -> PeerInfo {
let messenger_peer = self.messenger.peer_info();
let stream_addr = self.stream_transport.address();
if stream_addr.as_bytes().is_empty()
|| stream_addr
.available_transports()
.map(|v| v.is_empty())
.unwrap_or(true)
{
return messenger_peer;
}
let mut builder = crate::transports::address::WorkerAddressBuilder::new();
if let Err(e) = builder.merge(messenger_peer.worker_address()) {
tracing::warn!(
instance_id = %messenger_peer.instance_id(),
error = %e,
"peer_info: failed to merge messenger WorkerAddress into builder; \
falling back to messenger-only PeerInfo (streaming peers will not \
see this worker's streaming endpoint)"
);
return messenger_peer;
}
if let Err(e) = builder.merge(&stream_addr) {
tracing::warn!(
instance_id = %messenger_peer.instance_id(),
streaming_key = %self.stream_transport.key(),
error = %e,
"peer_info: failed to merge streaming WorkerAddress into builder; \
falling back to messenger-only PeerInfo (likely a key collision \
with a messenger transport key)"
);
return messenger_peer;
}
match builder.build() {
Ok(merged) => PeerInfo::new(messenger_peer.instance_id(), merged),
Err(e) => {
tracing::warn!(
instance_id = %messenger_peer.instance_id(),
error = %e,
"peer_info: WorkerAddressBuilder::build() failed; falling back \
to messenger-only PeerInfo"
);
messenger_peer
}
}
}
pub fn events(&self) -> &Arc<VeloEvents> {
self.messenger.events()
}
pub fn event_manager(&self) -> EventManager {
self.messenger.event_manager()
}
pub fn am_send(&self, handler: &str) -> Result<AmSendBuilder> {
self.messenger.am_send(handler)
}
pub fn am_sync(&self, handler: &str) -> Result<AmSyncBuilder> {
self.messenger.am_sync(handler)
}
pub fn unary(&self, handler: &str) -> Result<UnaryBuilder> {
self.messenger.unary(handler)
}
pub fn typed_unary<R: serde::de::DeserializeOwned + Send + 'static>(
&self,
handler: &str,
) -> Result<TypedUnaryBuilder<R>> {
self.messenger.typed_unary(handler)
}
pub fn register_handler(&self, handler: Handler) -> Result<()> {
self.messenger.register_handler(handler)
}
pub fn register_peer(&self, peer_info: PeerInfo) -> Result<()> {
let stream_key = self.stream_transport.key();
match peer_info.worker_address().get_entry(stream_key.as_str()) {
Ok(Some(_)) => {
self.stream_transport.register(&peer_info)?;
}
Ok(None) => {
tracing::debug!(
peer = %peer_info.worker_id(),
streaming_key = %stream_key,
"streaming transport register: peer has no matching streaming endpoint"
);
}
Err(e) => {
return Err(anyhow::anyhow!(
"decoding peer WorkerAddress for streaming key '{}': {e}",
stream_key
));
}
}
self.messenger.register_peer(peer_info)
}
pub async fn discover_and_register_peer(&self, instance_id: InstanceId) -> Result<()> {
let discovery = self.messenger.discovery().ok_or_else(|| {
anyhow::anyhow!(
"No discovery backend configured. Cannot discover instance {}",
instance_id
)
})?;
let peer_info = discovery.discover_by_instance_id(instance_id).await?;
self.register_peer(peer_info)
}
pub fn has_event_subscriber(&self, handle: EventHandle, subscriber: InstanceId) -> bool {
self.messenger.has_event_subscriber(handle, subscriber)
}
pub async fn available_handlers(&self, instance_id: InstanceId) -> Result<Vec<String>> {
self.messenger.available_handlers(instance_id).await
}
pub async fn refresh_handlers(&self, instance_id: InstanceId) -> Result<()> {
self.messenger.refresh_handlers(instance_id).await
}
pub async fn wait_for_handler(
&self,
instance_id: InstanceId,
handler_name: &str,
) -> Result<()> {
self.messenger
.wait_for_handler(instance_id, handler_name)
.await
}
pub fn list_local_handlers(&self) -> Vec<String> {
self.messenger.list_local_handlers()
}
pub fn runtime(&self) -> &tokio::runtime::Handle {
self.messenger.runtime()
}
pub fn tracker(&self) -> &tokio_util::task::TaskTracker {
self.messenger.tracker()
}
pub fn create_anchor<T>(&self) -> StreamAnchor<T> {
self.anchor_manager.create_anchor::<T>()
}
pub async fn attach_anchor<T: serde::Serialize>(
&self,
handle: StreamAnchorHandle,
) -> Result<StreamSender<T>, AttachError> {
self.anchor_manager.attach_stream_anchor::<T>(handle).await
}
pub fn anchor_manager(&self) -> &crate::streaming::AnchorManager {
&self.anchor_manager
}
pub fn create_mpsc_anchor<T>(&self) -> streaming::mpsc::MpscStreamAnchor<T> {
self.anchor_manager.create_mpsc_anchor::<T>()
}
pub fn create_mpsc_anchor_with_config<T>(
&self,
config: streaming::mpsc::MpscAnchorConfig,
) -> streaming::mpsc::MpscStreamAnchor<T> {
self.anchor_manager
.create_mpsc_anchor_with_config::<T>(config)
}
pub async fn attach_mpsc_anchor<T: serde::Serialize>(
&self,
handle: StreamAnchorHandle,
) -> Result<streaming::mpsc::MpscStreamSender<T>, AttachError> {
self.anchor_manager
.attach_mpsc_stream_anchor::<T>(handle)
.await
}
pub fn register_data(&self, data: bytes::Bytes) -> DataHandle {
self.rendezvous_manager.register_data(data)
}
pub fn register_data_with(&self, data: bytes::Bytes, opts: RegisterOptions) -> DataHandle {
self.rendezvous_manager.register_data_with(data, opts)
}
pub async fn register_data_pinned(&self, data: &[u8]) -> DataHandle {
self.rendezvous_manager.register_data_pinned(data).await
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub fn register_data_in_region(
&self,
guard: &RegionGuard,
range: std::ops::Range<u64>,
) -> Result<DataHandle, RdmaError> {
self.rendezvous_manager
.register_data_in_region(guard, range)
}
pub async fn metadata(&self, handle: DataHandle) -> Result<DataMetadata> {
self.rendezvous_manager.metadata(handle).await
}
pub async fn get(&self, handle: DataHandle) -> Result<(bytes::Bytes, u64)> {
self.rendezvous_manager.get(handle).await
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub async fn get_pinned(&self, handle: DataHandle) -> Result<(PinnedBuf, u64)> {
self.rendezvous_manager.get_pinned(handle).await
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub async fn alloc_pinned_writer(&self, len: usize) -> Result<PinnedWriter, RdmaError> {
self.rendezvous_manager.alloc_pinned_writer(len).await
}
pub async fn get_into(
&self,
handle: DataHandle,
dest: &mut impl RendezvousWrite,
) -> Result<u64> {
self.rendezvous_manager.get_into(handle, dest).await
}
pub async fn ref_handle(&self, handle: DataHandle) -> Result<()> {
self.rendezvous_manager.ref_handle(handle).await
}
pub async fn detach(&self, handle: DataHandle, lease_id: u64) -> Result<()> {
self.rendezvous_manager.detach(handle, lease_id).await
}
pub async fn release(&self, handle: DataHandle, lease_id: u64) -> Result<()> {
self.rendezvous_manager.release(handle, lease_id).await
}
#[cfg(feature = "test-helpers")]
pub fn primary_transport_key(&self, instance: InstanceId) -> Option<String> {
self.messenger
.backend()
.primary_transport_key(instance)
.map(|key| key.as_str().to_string())
}
pub fn rendezvous_manager(&self) -> &crate::rendezvous::RendezvousManager {
&self.rendezvous_manager
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub async unsafe fn register_external_memory(
&self,
ptr: std::ptr::NonNull<u8>,
len: usize,
) -> Result<RegionGuard, RdmaError> {
let rdma = self.rdma_registry()?;
unsafe { rdma.register_external(ptr, len) }.await
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub async fn register_owned(
&self,
buf: Box<[u8]>,
) -> Result<RegionGuard, crate::rendezvous::rdma::RegisterOwnedError> {
let Some(rdma) = self.rdma.as_ref() else {
return Err(crate::rendezvous::rdma::RegisterOwnedError {
buffer: Some(buf),
cause: RdmaError::NotConfigured,
});
};
rdma.register_owned(buf).await
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub fn rdma_registered_bytes(&self) -> u64 {
self.rdma
.as_ref()
.map(|r| r.registered_bytes())
.unwrap_or(0)
}
#[cfg(all(target_os = "linux", feature = "ucx", feature = "test-helpers"))]
pub fn rdma_in_flight_transfers(&self) -> usize {
self.rdma
.as_ref()
.map(|r| r.in_flight_transfers())
.unwrap_or(0)
}
#[cfg(all(target_os = "linux", feature = "ucx", test))]
pub(crate) fn rdma(&self) -> Option<&Arc<crate::rendezvous::rdma::RdmaRegistry>> {
self.rdma.as_ref()
}
#[cfg(all(target_os = "linux", feature = "ucx"))]
pub(crate) fn rdma_registry(
&self,
) -> Result<&Arc<crate::rendezvous::rdma::RdmaRegistry>, RdmaError> {
self.rdma.as_ref().ok_or(RdmaError::NotConfigured)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(all(target_os = "linux", feature = "ucx"))]
#[test]
fn the_rdma_kill_switch_reads_only_affirmatives() {
for on in [
"1", "true", "TRUE", "True", "yes", "YES", "on", "ON", " 1 ", "\ttrue\n",
] {
assert!(
rdma_rendezvous_disabled(Some(on)),
"{on:?} should switch the rendezvous RDMA path off"
);
}
for off in [
"0", "false", "no", "off", "", " ", "2", "disable", "ture", "1 1",
] {
assert!(
!rdma_rendezvous_disabled(Some(off)),
"{off:?} must not switch the rendezvous RDMA path off"
);
}
assert!(
!rdma_rendezvous_disabled(None),
"an unset variable must leave the path enabled"
);
}
#[test]
fn test_stream_config_double_call_error() {
let builder = Velo::builder();
let builder = builder
.stream_config(StreamConfig::Tcp(None))
.expect("first stream_config should succeed");
let result = builder.stream_config(StreamConfig::Tcp(None));
assert!(
result.is_err(),
"second stream_config call should return Err"
);
let err = result.err().unwrap();
assert!(
err.to_string().contains("more than once") || err.to_string().contains("one streaming"),
"error message should indicate double-call: {}",
err
);
}
#[test]
fn velo_has_anchor_manager_accessor() {
let _: fn(&Velo) -> &crate::streaming::AnchorManager = Velo::anchor_manager;
}
#[test]
fn velo_create_anchor_signature() {
let _: fn(&Velo) -> crate::streaming::StreamAnchor<String> = Velo::create_anchor::<String>;
}
#[tokio::test]
async fn velo_attach_anchor_type_checks() {
let transport = {
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
Arc::new(
crate::transports::tcp::TcpTransportBuilder::new()
.from_listener(listener)
.unwrap()
.build()
.unwrap(),
)
};
let velo = Velo::builder()
.add_transport(transport)
.build()
.await
.unwrap();
let _am: &crate::streaming::AnchorManager = velo.anchor_manager();
let anchor: crate::streaming::StreamAnchor<String> = velo.create_anchor::<String>();
let handle = anchor.handle();
let result: Result<crate::streaming::StreamSender<String>, crate::streaming::AttachError> =
velo.attach_anchor::<String>(handle).await;
let _sender = result.expect("local attach should succeed");
}
}