Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 16 additions & 4 deletions crates/ember-server/src/grpc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<crate::pubsub::PubMessage>)> = Vec::new();
let mut pattern_rxs: Vec<(String, tokio::sync::broadcast::Receiver<crate::pubsub::PubMessage>)> = Vec::new();
let mut channel_rxs: Vec<(
String,
tokio::sync::broadcast::Receiver<crate::pubsub::PubMessage>,
)> = Vec::new();
let mut pattern_rxs: Vec<(
String,
tokio::sync::broadcast::Receiver<crate::pubsub::PubMessage>,
)> = Vec::new();

for ch in &req.channels {
channel_rxs.push((ch.clone(), pubsub.subscribe(ch)));
Expand Down Expand Up @@ -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<crate::pubsub::PubMessage>)],
rxs: &mut [(
String,
tokio::sync::broadcast::Receiver<crate::pubsub::PubMessage>,
)],
) -> Option<SubscribeEvent> {
if rxs.is_empty() {
// no channel subscriptions — park forever so pattern branch can drive
Expand Down Expand Up @@ -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<crate::pubsub::PubMessage>)],
rxs: &mut [(
String,
tokio::sync::broadcast::Receiver<crate::pubsub::PubMessage>,
)],
) -> Option<SubscribeEvent> {
if rxs.is_empty() {
std::future::pending::<()>().await;
Expand Down
16 changes: 12 additions & 4 deletions crates/ember-server/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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())
Expand Down Expand Up @@ -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())
Expand Down