use std::sync::Arc;
use anyhow::Result;
pub use backend::Transport;
pub use velo_observability::VeloMetrics;
pub use velo_transports as backend;
pub use velo_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 velo_common::{InstanceId, PeerInfo, WorkerAddress, WorkerId};
pub use velo_events::{
Event, EventAwaiter, EventBackend, EventHandle, EventManager, EventPoison, EventStatus,
};
pub use velo_discovery as discovery;
pub use velo_streaming::{
AnchorManager, AttachError, SendError, StreamAnchor, StreamAnchorHandle, StreamController,
StreamError, StreamFrame, StreamSender,
};
pub use velo_queue as queue;
pub use velo_rendezvous::{
DataHandle, DataMetadata, RegisterOptions, RendezvousManager, RendezvousWrite, StageMode,
};
#[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<velo_streaming::AnchorManager>,
rendezvous_manager: Arc<velo_rendezvous::RendezvousManager>,
}
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 velo_transport = Arc::new(velo_streaming::VeloFrameTransport::new(
Arc::clone(&messenger),
worker_id,
self.metrics.clone(),
)?);
let (default_transport, transport_registry) = match self.stream_config {
Some(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_transport = Arc::new(velo_streaming::TcpFrameTransport::new(bind_addr));
let mut registry = std::collections::HashMap::new();
registry.insert(
"tcp".to_string(),
Arc::clone(&tcp_transport) as Arc<dyn velo_streaming::FrameTransport>,
);
registry.insert(
"velo".to_string(),
velo_transport.clone() as Arc<dyn velo_streaming::FrameTransport>,
);
(
tcp_transport as Arc<dyn velo_streaming::FrameTransport>,
Arc::new(registry),
)
}
#[cfg(feature = "grpc")]
Some(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_transport = Arc::new(
velo_streaming::GrpcFrameTransport::new(bind_addr)
.await
.map_err(|e| anyhow::anyhow!("Failed to start gRPC transport: {}", e))?,
);
let mut registry = std::collections::HashMap::new();
registry.insert(
"grpc".to_string(),
Arc::clone(&grpc_transport) as Arc<dyn velo_streaming::FrameTransport>,
);
registry.insert(
"velo".to_string(),
velo_transport.clone() as Arc<dyn velo_streaming::FrameTransport>,
);
(
grpc_transport as Arc<dyn velo_streaming::FrameTransport>,
Arc::new(registry),
)
}
None => (
velo_transport as Arc<dyn velo_streaming::FrameTransport>,
Arc::new(std::collections::HashMap::new()),
),
};
let anchor_manager = Arc::new(
velo_streaming::AnchorManagerBuilder::default()
.worker_id(worker_id)
.transport(default_transport)
.transport_registry(transport_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) => velo_rendezvous::RendezvousManager::with_metrics(worker_id, Arc::clone(m)),
None => velo_rendezvous::RendezvousManager::new(worker_id),
});
rendezvous_manager.register_handlers(Arc::clone(&messenger))?;
let stager = Arc::new(velo_rendezvous::RendezvousStager::new(Arc::clone(
&rendezvous_manager,
)));
let resolver = Arc::new(velo_rendezvous::RendezvousResolver::new(Arc::clone(
&rendezvous_manager,
)));
messenger.set_large_payload_support(stager, resolver);
Ok(Arc::new(Velo {
messenger,
anchor_manager,
rendezvous_manager,
}))
}
}
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 {
self.messenger.peer_info()
}
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<()> {
self.messenger.register_peer(peer_info)
}
pub async fn discover_and_register_peer(&self, instance_id: InstanceId) -> Result<()> {
self.messenger.discover_and_register_peer(instance_id).await
}
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) -> &velo_streaming::AnchorManager {
&self.anchor_manager
}
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) -> &velo_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) -> &velo_streaming::AnchorManager = Velo::anchor_manager;
}
#[test]
fn velo_create_anchor_signature() {
let _: fn(&Velo) -> velo_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(
velo_transports::tcp::TcpTransportBuilder::new()
.from_listener(listener)
.unwrap()
.build()
.unwrap(),
)
};
let velo = Velo::builder()
.add_transport(transport)
.build()
.await
.unwrap();
let _am: &velo_streaming::AnchorManager = velo.anchor_manager();
let anchor: velo_streaming::StreamAnchor<String> = velo.create_anchor::<String>();
let handle = anchor.handle();
let result: Result<velo_streaming::StreamSender<String>, velo_streaming::AttachError> =
velo.attach_anchor::<String>(handle).await;
let _sender = result.expect("local attach should succeed");
}
}