diff --git a/async-nats/src/jetstream/consumer/push.rs b/async-nats/src/jetstream/consumer/push.rs index 327807170..d2b37f696 100644 --- a/async-nats/src/jetstream/consumer/push.rs +++ b/async-nats/src/jetstream/consumer/push.rs @@ -42,7 +42,7 @@ use std::{sync::atomic::Ordering, time::Duration}; #[cfg(feature = "server_2_11")] use time::{serde::rfc3339, OffsetDateTime}; use tokio::{sync::oneshot::error::TryRecvError, task::JoinHandle}; -use tracing::{debug, trace}; +use tracing::{debug, error, trace}; const ORDERED_IDLE_HEARTBEAT: Duration = Duration::from_secs(5); @@ -153,10 +153,14 @@ impl futures_util::Stream for Messages { // TODO store pending_publish as a future and return errors from it let client = self.context.client.clone(); tokio::task::spawn(async move { - client - .publish(subject, Bytes::from_static(b"")) - .await - .unwrap(); + if let Err(err) = + client.publish(subject, Bytes::from_static(b"")).await + { + error!( + "failed to respond to idle heartbeat: {}", + err + ); + } }); } diff --git a/async-nats/src/service/mod.rs b/async-nats/src/service/mod.rs index 418b55bfa..dba0465f1 100644 --- a/async-nats/src/service/mod.rs +++ b/async-nats/src/service/mod.rs @@ -380,6 +380,10 @@ impl Service { loop { tokio::select! { Some(ping) = pings.next() => { + let Some(reply) = ping.reply else { + debug!("ignoring PING request without a reply subject"); + continue; + }; let pong = serde_json::to_vec(&PingResponse{ kind: "io.nats.micro.v1.ping_response".to_string(), name: info.name.clone(), @@ -387,9 +391,13 @@ impl Service { version: info.version.clone(), metadata: info.metadata.clone(), })?; - client.publish(ping.reply.unwrap(), pong.into()).await?; + client.publish(reply, pong.into()).await?; }, Some(info_request) = infos.next() => { + let Some(reply) = info_request.reply else { + debug!("ignoring INFO request without a reply subject"); + continue; + }; let info = info.clone(); let endpoints: Vec = { @@ -407,9 +415,13 @@ impl Service { ..info }; let info_json = serde_json::to_vec(&info).map(Bytes::from)?; - client.publish(info_request.reply.unwrap(), info_json.clone()).await?; + client.publish(reply, info_json.clone()).await?; }, Some(stats_request) = stats.next() => { + let Some(reply) = stats_request.reply else { + debug!("ignoring STATS request without a reply subject"); + continue; + }; if let Some(stats_callback) = stats_callback.as_mut() { let mut endpoint_stats_locked = endpoints_state.lock().unwrap(); for (key, value) in &mut endpoint_stats_locked.endpoints { @@ -425,7 +437,7 @@ impl Service { started, endpoints: endpoints_state.lock().unwrap().endpoints.values().cloned().map(Into::into).collect(), })?; - client.publish(stats_request.reply.unwrap(), stats.into()).await?; + client.publish(reply, stats.into()).await?; }, else => break, }