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::{InstanceId, PeerInfo, Transport, WorkerAddress, WorkerId};
pub use velo_ext as ext;
pub use crate::messenger::{
AmHandlerBuilder, AmSendBuilder, AmSyncBuilder, AsyncExecutor, Context, DispatchMode, Handler,
HandlerExecutor, Messenger, MessengerBuilder, 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,
};
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>,
}
pub struct VeloBuilder {
inner: MessengerBuilder,
stream_config: Option<StreamConfig>,
metrics: Option<Arc<VeloMetrics>>,
}
impl VeloBuilder {
pub fn new() -> Self {
Self {
inner: MessengerBuilder::new(),
stream_config: None,
metrics: None,
}
}
pub fn add_transport(mut self, transport: Arc<dyn Transport>) -> Self {
self.inner = self.inner.add_transport(transport);
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 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 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))?,
);
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);
Ok(Arc::new(Velo {
messenger,
anchor_manager,
rendezvous_manager,
stream_transport,
}))
}
}
impl Default for VeloBuilder {
fn default() -> Self {
Self::new()
}
}
impl Velo {
pub fn builder() -> VeloBuilder {
VeloBuilder::new()
}
pub fn messenger(&self) -> &Arc<Messenger> {
&self.messenger
}
pub fn instance_id(&self) -> InstanceId {
self.messenger.instance_id()
}
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 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
}
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
}
pub fn rendezvous_manager(&self) -> &crate::rendezvous::RendezvousManager {
&self.rendezvous_manager
}
}
#[cfg(test)]
mod tests {
use super::*;
#[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");
}
}