From b2e366c739f2939abf914115fc151b1a2e788150 Mon Sep 17 00:00:00 2001 From: Ximon Eighteen <3304436+ximon18@users.noreply.github.com> Date: Mon, 3 Aug 2026 13:13:46 +0200 Subject: [PATCH] FIX: Abort response streaming if no more responses can be sent. This prevents serving of an AXFR continuing to attempt to push messages into the response even after the connection to the client has been closed or no more response messages can be enqueued. --- src/net/server/connection.rs | 23 +++++++++++++++-------- 1 file changed, 15 insertions(+), 8 deletions(-) diff --git a/src/net/server/connection.rs b/src/net/server/connection.rs index 674b5c92..b83cf580 100644 --- a/src/net/server/connection.rs +++ b/src/net/server/connection.rs @@ -1,5 +1,7 @@ //! Support for stream based connections. +use core::fmt::Display; use core::future::Future; +use core::net::SocketAddr; use core::ops::{ControlFlow, Deref}; use core::pin::Pin; use core::sync::atomic::{AtomicBool, Ordering}; @@ -7,9 +9,8 @@ use core::time::Duration; use alloc::boxed::Box; use alloc::sync::Arc; -use core::fmt::Display; -use core::net::SocketAddr; use std::io; +use std::sync::Mutex; use arc_swap::ArcSwap; use log::{Level, log_enabled}; @@ -1036,7 +1037,7 @@ where metrics: Arc, /// The status of the service invoker. - status: InvokerStatus, + status: Mutex, } impl @@ -1058,7 +1059,7 @@ where config, result_q_tx, metrics, - status: InvokerStatus::Normal, + status: Mutex::new(InvokerStatus::Normal), } } @@ -1096,11 +1097,15 @@ where error!( "Unable to queue message for sending: connection is shutting down." ); + *self.status.lock().unwrap() = InvokerStatus::Aborting; break; } Err(TrySendError::Full(unused_response)) => { - if matches!(self.status, InvokerStatus::InTransaction) { + if matches!( + *self.status.lock().unwrap(), + InvokerStatus::InTransaction + ) { // Wait until there is space in the message queue. tokio::task::yield_now().await; response = unused_response; @@ -1108,6 +1113,8 @@ where error!( "Unable to queue message for sending: queue is full." ); + *self.status.lock().unwrap() = + InvokerStatus::Aborting; break; } } @@ -1130,7 +1137,7 @@ where config: self.config.clone(), result_q_tx: self.result_q_tx.clone(), metrics: self.metrics.clone(), - status: InvokerStatus::Normal, + status: Mutex::new(InvokerStatus::Normal), } } } @@ -1146,11 +1153,11 @@ where Svc: Service + Clone, { fn status(&self) -> InvokerStatus { - self.status + *self.status.lock().unwrap() } fn set_status(&mut self, status: InvokerStatus) { - self.status = status; + *self.status.lock().unwrap() = status; } fn reconfigure(&self, idle_timeout: Option) {