diff --git a/crates/ember-server/src/grpc.rs b/crates/ember-server/src/grpc.rs index 8a1fc8ba..45772319 100644 --- a/crates/ember-server/src/grpc.rs +++ b/crates/ember-server/src/grpc.rs @@ -1845,8 +1845,14 @@ impl EmberCache for EmberService { let pubsub = Arc::clone(&self.pubsub); // collect all broadcast receivers - let mut channel_rxs: Vec<(String, tokio::sync::broadcast::Receiver)> = Vec::new(); - let mut pattern_rxs: Vec<(String, tokio::sync::broadcast::Receiver)> = Vec::new(); + let mut channel_rxs: Vec<( + String, + tokio::sync::broadcast::Receiver, + )> = Vec::new(); + let mut pattern_rxs: Vec<( + String, + tokio::sync::broadcast::Receiver, + )> = Vec::new(); for ch in &req.channels { channel_rxs.push((ch.clone(), pubsub.subscribe(ch))); @@ -2095,7 +2101,10 @@ async fn handle_pipeline_command( /// Waits for the next message on any channel subscription receiver. /// Returns None when all receivers are closed. async fn recv_any_channel( - rxs: &mut [(String, tokio::sync::broadcast::Receiver)], + rxs: &mut [( + String, + tokio::sync::broadcast::Receiver, + )], ) -> Option { if rxs.is_empty() { // no channel subscriptions — park forever so pattern branch can drive @@ -2129,7 +2138,10 @@ async fn recv_any_channel( /// Waits for the next message on any pattern subscription receiver. /// Returns None when all receivers are closed. async fn recv_any_pattern( - rxs: &mut [(String, tokio::sync::broadcast::Receiver)], + rxs: &mut [( + String, + tokio::sync::broadcast::Receiver, + )], ) -> Option { if rxs.is_empty() { std::future::pending::<()>().await; diff --git a/crates/ember-server/src/server.rs b/crates/ember-server/src/server.rs index b3a08e38..57e9fda5 100644 --- a/crates/ember-server/src/server.rs +++ b/crates/ember-server/src/server.rs @@ -131,8 +131,12 @@ pub async fn run( // spawn gRPC listener if configured #[cfg(feature = "grpc")] let _grpc_handle = if let Some(grpc_addr) = grpc_addr { - let svc = - crate::grpc::EmberService::new(engine.clone(), Arc::clone(&ctx), Arc::clone(&slow_log), Arc::clone(&pubsub)); + let svc = crate::grpc::EmberService::new( + engine.clone(), + Arc::clone(&ctx), + Arc::clone(&slow_log), + Arc::clone(&pubsub), + ); info!("gRPC listening on {grpc_addr}"); let server = tonic::transport::Server::builder() .add_service(svc.into_service()) @@ -396,8 +400,12 @@ pub async fn run_concurrent( // spawn gRPC listener if configured #[cfg(feature = "grpc")] let _grpc_handle = if let Some(grpc_addr) = grpc_addr { - let svc = - crate::grpc::EmberService::new(engine.clone(), Arc::clone(&ctx), Arc::clone(&slow_log), Arc::clone(&pubsub)); + let svc = crate::grpc::EmberService::new( + engine.clone(), + Arc::clone(&ctx), + Arc::clone(&slow_log), + Arc::clone(&pubsub), + ); info!("gRPC listening on {grpc_addr}"); let server = tonic::transport::Server::builder() .add_service(svc.into_service())