From ed58ed03258f6e5d4af3df130cd0aef4c8ae44ee Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Sun, 13 Sep 2026 22:29:56 -0600 Subject: [PATCH] Pipeline compatible article requests upstream --- benches/cache_miss_roundtrip_callgrind.rs | 82 +++-- src/proxy/lifecycle.rs | 22 +- src/session/backend.rs | 18 ++ src/session/handlers/command_execution.rs | 97 +++++- src/session/handlers/per_command.rs | 368 ++++++++++++++++++++-- tests/rfc3977/body_workflow.rs | 180 +++++++++++ 6 files changed, 704 insertions(+), 63 deletions(-) diff --git a/benches/cache_miss_roundtrip_callgrind.rs b/benches/cache_miss_roundtrip_callgrind.rs index 9cef278f..64c417b1 100644 --- a/benches/cache_miss_roundtrip_callgrind.rs +++ b/benches/cache_miss_roundtrip_callgrind.rs @@ -7,16 +7,11 @@ //! //! Run with: `cargo bench --bench cache_miss_roundtrip_callgrind` -macro_rules! supported { - ($($item:item)*) => { - $( - #[cfg(all(target_os = "linux", any(target_arch = "x86_64", target_arch = "aarch64")))] - $item - )* - }; -} - -supported! { +#[cfg(all( + target_os = "linux", + any(target_arch = "x86_64", target_arch = "aarch64") +))] +mod benchmarks { use gungraun::{ Callgrind, LibraryBenchmarkConfig, library_benchmark, library_benchmark_group, main, }; @@ -49,10 +44,10 @@ supported! { Config { servers: vec![ Server::builder("127.0.0.1", Port::try_new(backend_port).unwrap()) - .name("bench-backend") - .max_connections(MaxConnections::try_new(4).unwrap()) - .build() - .unwrap(), + .name("bench-backend") + .max_connections(MaxConnections::try_new(4).unwrap()) + .build() + .unwrap(), ], cache, ..Default::default() @@ -191,8 +186,36 @@ supported! { } async fn article_roundtrip(&mut self) -> usize { - self.stream.write_all(b"ARTICLE \r\n").await.unwrap(); - read_exact_response_into(&mut self.stream, &mut self.response_buffer, self.response_len).await + self.stream + .write_all(b"ARTICLE \r\n") + .await + .unwrap(); + read_exact_response_into( + &mut self.stream, + &mut self.response_buffer, + self.response_len, + ) + .await + } + + async fn article_pair_roundtrip(&mut self) -> usize { + self.stream + .write_all(b"ARTICLE \r\nARTICLE \r\n") + .await + .unwrap(); + let first = read_exact_response_into( + &mut self.stream, + &mut self.response_buffer, + self.response_len, + ) + .await; + let second = read_exact_response_into( + &mut self.stream, + &mut self.response_buffer, + self.response_len, + ) + .await; + first + second } } @@ -207,10 +230,14 @@ supported! { } } - async fn read_exact_response_into(stream: &mut TcpStream, buffer: &mut [u8], expected: usize) -> usize { + async fn read_exact_response_into( + stream: &mut TcpStream, + buffer: &mut [u8], + expected: usize, + ) -> usize { let mut total = 0usize; while total < expected { - let n = stream.read(&mut buffer[total..]).await.unwrap(); + let n = stream.read(&mut buffer[total..expected]).await.unwrap(); assert_ne!(n, 0, "proxy closed during benchmark response"); total += n; } @@ -249,6 +276,12 @@ supported! { black_box(harness.rt.block_on(harness.client.article_roundtrip())) } + #[library_benchmark] + #[bench::article_64k(args = (ARTICLE_64K), setup = setup_no_configured_cache_roundtrip)] + fn run_no_configured_cache_pair_roundtrip(mut harness: BenchHarness) -> usize { + black_box(harness.rt.block_on(harness.client.article_pair_roundtrip())) + } + #[library_benchmark] #[bench::article_64k(args = (ARTICLE_64K), setup = setup_metadata_only_cache_roundtrip)] #[bench::article_768k(args = (ARTICLE_768K), setup = setup_metadata_only_cache_roundtrip)] @@ -267,6 +300,7 @@ supported! { name = cache_miss_roundtrip; benchmarks = run_no_configured_cache_roundtrip, + run_no_configured_cache_pair_roundtrip, run_metadata_only_cache_roundtrip, run_direct_backend_roundtrip ); @@ -280,6 +314,18 @@ supported! { ); library_benchmark_groups = cache_miss_roundtrip ); + + pub(super) fn run() { + main(); + } +} + +#[cfg(all( + target_os = "linux", + any(target_arch = "x86_64", target_arch = "aarch64") +))] +fn main() { + benchmarks::run(); } #[cfg(not(all( diff --git a/src/proxy/lifecycle.rs b/src/proxy/lifecycle.rs index 305303fc..baae6cde 100644 --- a/src/proxy/lifecycle.rs +++ b/src/proxy/lifecycle.rs @@ -402,19 +402,23 @@ impl NntpProxy { /// /// This creates a session with the router, allowing commands from this client /// to be routed to different backends based on load balancing. - pub async fn handle_client_per_command_routing( + pub fn handle_client_per_command_routing( &self, client_stream: TcpStream, client_addr: ClientAddress, - ) -> Result<(), SessionError> { - // Check for stale pools before handling (lazy recreation after idle) - self.check_and_clear_stale_pools(); - self.increment_active_clients(); - - let result = Box::pin(self.handle_per_command_client(client_stream, client_addr)).await; + ) -> futures::future::BoxFuture<'_, Result<(), SessionError>> { + Box::pin(async move { + // Check for stale pools before handling (lazy recreation after idle) + self.check_and_clear_stale_pools(); + self.increment_active_clients(); + + let result = self + .handle_per_command_client(client_stream, client_addr) + .await; - self.decrement_active_clients(); - result + self.decrement_active_clients(); + result + }) } /// Handle a per-command routing session diff --git a/src/session/backend.rs b/src/session/backend.rs index a0587283..2466deaf 100644 --- a/src/session/backend.rs +++ b/src/session/backend.rs @@ -273,6 +273,24 @@ where read_until_backend_reply(conn, request, buffer).await } +/// Read a response for a request that was already written as part of an +/// upstream pipeline window. +pub(crate) async fn read_presend_request_classified( + conn: &mut C, + request: &RequestContext, + buffer: &mut PooledBuffer, +) -> Result +where + C: AsyncReadExt + Unpin, +{ + let n = buffer.read_from(conn).await?; + if n == 0 { + anyhow::bail!("Backend connection closed unexpectedly"); + } + + read_until_backend_reply(conn, request, buffer).await +} + pub(crate) async fn execute_request_classified_timed( conn: &mut C, request: &RequestContext, diff --git a/src/session/handlers/command_execution.rs b/src/session/handlers/command_execution.rs index 51078c5a..23148ff7 100644 --- a/src/session/handlers/command_execution.rs +++ b/src/session/handlers/command_execution.rs @@ -143,6 +143,11 @@ enum BackendReadAttemptError { Backend(anyhow::Error), } +pub(super) enum PresendResponseError { + Read(anyhow::Error), + Transfer(ResponseTransferError), +} + impl From for BackendReadAttemptError { fn from(err: SessionError) -> Self { Self::Backend(anyhow::Error::new(err)) @@ -821,7 +826,7 @@ impl ClientSession { )) } - fn handle_response_transfer_error( + pub(super) fn handle_response_transfer_error( &self, conn: crate::pool::ConnectionGuard, backend_id: BackendId, @@ -940,7 +945,7 @@ impl ClientSession { } } - async fn checkout_direct_backend_connection( + pub(super) async fn checkout_direct_backend_connection( &self, provider: &crate::pool::DeadpoolConnectionProvider, backend_id: crate::types::BackendId, @@ -1247,6 +1252,92 @@ impl ClientSession { Ok((response, buffer, timings)) } + pub(super) async fn read_and_write_presend_response( + &self, + conn: &mut crate::pool::ConnectionGuard, + client_write: &mut W, + backend: &ArticleBackend, + request: &mut RequestContext, + availability: &mut crate::cache::ArticleAvailability, + backend_to_client_bytes: &mut BackendToClientBytes, + ) -> Result + where + W: AsyncWrite + Unpin, + { + let backend_id = backend.backend_id(); + let mut buffer = self.buffer_pool.acquire(); + let read = + backend::read_presend_request_classified(conn.stream_mut(), request, &mut buffer) + .await + .map_err(PresendResponseError::Read)?; + let Some(status_code) = read.status_code() else { + read.log_warnings(&buffer, self.client_addr, backend_id); + return Err(PresendResponseError::Read(anyhow::anyhow!( + "backend returned an invalid response to a presend request" + ))); + }; + + if let Ok(missing) = AuthoritativeArticleMissing::from_status_code(*backend, status_code) { + self.record_authoritative_article_missing(&missing, availability); + if let Some(article_request) = + crate::command::CommandHandler::article_lookup_request(request) + { + self.cache + .record_availability_missing( + article_request.message_id(), + missing.availability_slot(), + ) + .await; + } + let completion = backend::observe_response( + request, + &mut buffer, + conn.stream_mut(), + &self.buffer_pool, + backend_id, + ) + .await + .map_err(PresendResponseError::Transfer)?; + self.send_430_to_client(client_write, backend_to_client_bytes) + .await + .map_err(|error| { + PresendResponseError::Transfer(classify_response_write_err(error)) + })?; + client_write.flush().await.map_err(|error| { + PresendResponseError::Transfer(classify_response_write_err(error)) + })?; + return Ok(completion); + } + + let params = ResponseWriteParams { + request, + article_request: crate::command::CommandHandler::article_lookup_request(request), + status_code, + }; + let (bytes_written, completion) = self + .write_response_to_client(conn.stream_mut(), client_write, backend, buffer, params) + .await + .map_err(PresendResponseError::Transfer)?; + client_write + .flush() + .await + .map_err(|error| PresendResponseError::Transfer(classify_response_write_err(error)))?; + self.record_response_metrics( + backend_id, + request, + status_code, + request.request_wire_len().as_u64(), + bytes_written, + ); + let response = RequestResponseMetadata::new( + status_code, + usize::try_from(bytes_written).unwrap_or(usize::MAX).into(), + ); + request.record_backend_response(backend_id, response); + *backend_to_client_bytes = backend_to_client_bytes.add(response.wire_len().get()); + Ok(completion) + } + #[inline] fn should_use_stat_missing_probe( provider: &crate::pool::DeadpoolConnectionProvider, @@ -1493,7 +1584,7 @@ impl ClientSession { &self, client_write: &mut W, backend_to_client_bytes: &mut BackendToClientBytes, - ) -> Result<()> + ) -> std::io::Result<()> where W: AsyncWrite + Unpin, { diff --git a/src/session/handlers/per_command.rs b/src/session/handlers/per_command.rs index bd9f6027..019044ac 100644 --- a/src/session/handlers/per_command.rs +++ b/src/session/handlers/per_command.rs @@ -11,12 +11,13 @@ use super::BackendLease; use crate::protocol::{ - AUTH_REQUIRED_FOR_COMMAND, RequestContext, RequestKind, RequestResponseMetadata, StatusCode, - codes, + AUTH_REQUIRED_FOR_COMMAND, RequestCacheStatus, RequestContext, RequestKind, + RequestResponseMetadata, StatusCode, codes, }; use crate::session::common; use crate::session::{ClientAuthState, ClientSession, connection}; use anyhow::Result; +use smallvec::SmallVec; use std::sync::Arc; use tokio::io::{AsyncWriteExt, BufReader}; use tokio::net::TcpStream; @@ -153,6 +154,23 @@ struct PerCommandLoopState { auth_access: AuthenticationAccess, } +enum PresendDecision { + Backend { + backend: crate::router::ArticleBackend, + availability: crate::cache::ArticleAvailability, + guard: Option, + }, + Missing, +} + +struct PipelineCommandParams<'a> { + router: &'a Arc, + client_writer: &'a mut crate::session::ClientWriter, + state: &'a mut PerCommandLoopState, + batch: &'a mut crate::session::handlers::pipeline::RequestBatch, + backend_connection: &'a mut BatchBackendConnection, +} + impl PerCommandLoopState { const fn new(auth_access: AuthenticationAccess) -> Self { Self { @@ -601,55 +619,339 @@ impl ClientSession { batch: &mut crate::session::handlers::pipeline::RequestBatch, backend_connection: &mut BatchBackendConnection, ) -> Result<(), SessionError> { + state.auth_access = self.authentication_access(state.auth_access); + if self + .try_process_presend_window(router, client_writer, state, batch, backend_connection) + .await? + { + return Ok(()); + } + let batch_size = batch.len(); for i in 0..batch_size { - let request = batch.context(i).request(); - debug!( - "Client {} received {} request bytes: kind={:?}, verb={:?}", - self.client_addr, - request.request_wire_len().get(), - request.kind(), - request.verb() - ); + self.process_pipelineable_command( + i, + true, + PipelineCommandParams { + router, + client_writer, + state, + batch, + backend_connection, + }, + ) + .await?; + } + + Ok(()) + } + async fn process_pipelineable_command( + &self, + i: usize, + account_request: bool, + params: PipelineCommandParams<'_>, + ) -> Result<(), SessionError> { + let PipelineCommandParams { + router, + client_writer, + state, + batch, + backend_connection, + } = params; + let request = batch.context(i).request(); + debug!( + "Client {} received {} request bytes: kind={:?}, verb={:?}", + self.client_addr, + request.request_wire_len().get(), + request.kind(), + request.verb() + ); + + if account_request { state.client_to_backend_bytes = state .client_to_backend_bytes .add(request.request_wire_len().get()); - state.auth_access = self.authentication_access(state.auth_access); + } + state.auth_access = self.authentication_access(state.auth_access); + + let command_result = batch + .with_context_mut(i, |mut request| async { + let result = self + .process_single_command(CommandExecutionParams { + request: &mut request, + auth_access: state.auth_access, + router, + client_writer, + backend_connection: backend_connection.slot(), + auth_username: &mut state.auth_username, + client_to_backend_bytes: state.client_to_backend_bytes, + backend_to_client_bytes: &mut state.backend_to_client_bytes, + }) + .await; + (request, result) + }) + .await + .map_err(|_| { + SessionError::Backend(anyhow::anyhow!( + "pipelineable request lost its pipelineability during execution" + )) + })?; + match command_result? { + SingleCommandResult::Continue => {} + SingleCommandResult::Quit => { + return Err(SessionError::Backend(anyhow::anyhow!( + "pipelineable request unexpectedly terminated the session" + ))); + } + } + + Ok(()) + } - let command_result = batch + async fn try_process_presend_window( + &self, + router: &Arc, + client_writer: &mut crate::session::ClientWriter, + state: &mut PerCommandLoopState, + batch: &mut crate::session::handlers::pipeline::RequestBatch, + backend_connection: &mut BatchBackendConnection, + ) -> Result { + if batch.len() < 2 + || router.backend_count() != 1 + || self.cache.stores_payload_responses() + || !state.auth_access.can_access_backend() + || (0..batch.len()).any(|i| { + let request = batch.context(i).request(); + !matches!( + request.kind(), + RequestKind::Article | RequestKind::Body | RequestKind::Head + ) || !matches!( + CommandHandler::classify_request( + request, + state.auth_access, + self.mode_state.routing_mode(), + ), + CommandPlan::Forward + ) + }) + { + return Ok(false); + } + + let Some(backend_id) = router + .tiers() + .flat_map(|tier| router.backend_ids_in_tier(tier)) + .next() + else { + return Ok(false); + }; + let Some(provider) = router.backend_provider(backend_id) else { + return Ok(false); + }; + if provider.stat_missing_enabled() { + return Ok(false); + } + + let mut decisions = SmallVec::<[PresendDecision; 16]>::new(); + for i in 0..batch.len() { + let availability = batch .with_context_mut(i, |mut request| async { - let result = self - .process_single_command(CommandExecutionParams { - request: &mut request, - auth_access: state.auth_access, - router, - client_writer, - backend_connection: backend_connection.slot(), - auth_username: &mut state.auth_username, - client_to_backend_bytes: state.client_to_backend_bytes, - backend_to_client_bytes: &mut state.backend_to_client_bytes, - }) - .await; - (request, result) + let cached = match request.message_id() { + Some(message_id) => self.cache.get_request_message_id(message_id).await, + None => None, + }; + let availability = if let Some(cached) = cached { + let availability = cached.to_availability(); + request.record_cache_entry_metadata( + cached.request_cache_metadata(&availability), + ); + request.record_cache_status(RequestCacheStatus::PartialHit); + availability + } else { + request.record_cache_status(RequestCacheStatus::Miss); + crate::cache::ArticleAvailability::new() + }; + (request, availability) }) .await .map_err(|_| { SessionError::Backend(anyhow::anyhow!( - "pipelineable request lost its pipelineability during execution" + "presend request lost its pipelineability during preparation" )) })?; - match command_result? { - SingleCommandResult::Continue => {} - SingleCommandResult::Quit => { - return Err(SessionError::Backend(anyhow::anyhow!( - "pipelineable request unexpectedly terminated the session" - ))); + let Some(availability_slot) = router.availability_slot(backend_id) else { + return Ok(false); + }; + let Some(backend) = crate::router::ArticleBackend::from_availability_slot( + backend_id, + availability_slot, + &availability, + ) else { + decisions.push(PresendDecision::Missing); + continue; + }; + decisions.push(PresendDecision::Backend { + backend, + availability, + guard: Some(BackendSelector::guard_for_manual_backend( + router.clone(), + backend_id, + )), + }); + } + + let Some(first_backend_index) = decisions + .iter() + .position(|decision| matches!(decision, PresendDecision::Backend { .. })) + else { + for _ in 0..batch.len() { + self.send_430_to_client( + client_writer.get_mut(), + &mut state.backend_to_client_bytes, + ) + .await?; + } + client_writer.get_mut().flush().await?; + return Ok(true); + }; + let first_request = batch.context(first_backend_index).request(); + let mut conn = match self + .checkout_direct_backend_connection( + provider, + backend_id, + first_request, + backend_connection.slot(), + ) + .await + { + Ok(conn) => conn, + Err(_) => return Ok(false), + }; + + for (i, decision) in decisions.iter().enumerate() { + if matches!(decision, PresendDecision::Backend { .. }) { + let request = batch.context(i).request(); + self.metrics.record_command(backend_id); + self.metrics.user_command(self.username()); + if request.write_wire_to(conn.stream_mut()).await.is_err() { + conn.fail_backend(); + return Ok(false); } } } + if conn.stream_mut().flush().await.is_err() { + conn.fail_backend(); + return Ok(false); + } + for i in 0..batch.len() { + state.client_to_backend_bytes = state + .client_to_backend_bytes + .add(batch.context(i).request().request_wire_len().get()); + } - Ok(()) + let mut last_completion = None; + for (i, decision) in decisions.iter_mut().enumerate() { + match decision { + PresendDecision::Missing => { + self.send_430_to_client( + client_writer.get_mut(), + &mut state.backend_to_client_bytes, + ) + .await?; + } + PresendDecision::Backend { + backend, + availability, + guard, + } => { + let result = batch + .with_context_mut(i, |mut request| async { + let result = self + .read_and_write_presend_response( + &mut conn, + client_writer.get_mut(), + backend, + &mut request, + availability, + &mut state.backend_to_client_bytes, + ) + .await; + (request, result) + }) + .await + .map_err(|_| { + SessionError::Backend(anyhow::anyhow!( + "presend request lost its pipelineability while reading response" + )) + })?; + match result { + Ok(completion) => { + last_completion = Some(completion); + guard + .take() + .expect("presend backend decision must retain its command guard") + .complete(); + } + Err( + crate::session::handlers::command_execution::PresendResponseError::Read( + error, + ), + ) => { + debug!( + client = %self.client_addr, + backend = backend_id.as_index(), + error = %error, + "Presend response read failed; retrying remaining requests sequentially" + ); + conn.fail_backend(); + for decision in &mut decisions[i..] { + if let PresendDecision::Backend { guard, .. } = decision { + drop(guard.take()); + } + } + for retry_index in i..batch.len() { + self.process_pipelineable_command( + retry_index, + false, + PipelineCommandParams { + router, + client_writer, + state, + batch, + backend_connection, + }, + ) + .await?; + } + return Ok(true); + } + Err( + crate::session::handlers::command_execution::PresendResponseError::Transfer( + error, + ), + ) => { + let request = batch.context(i).request(); + return Err(self.handle_response_transfer_error( + conn, + backend_id, + request, + error, + )); + } + } + } + } + } + client_writer.get_mut().flush().await?; + + let completion = last_completion.expect("presend window issued at least one request"); + if conn.has_pending_bytes() { + conn.fail_client(); + } else { + *backend_connection.slot() = Some(BackendLease::new(backend_id, conn, completion)); + } + Ok(true) } async fn handle_trailing_command( diff --git a/tests/rfc3977/body_workflow.rs b/tests/rfc3977/body_workflow.rs index 4cd52194..1a5cb671 100644 --- a/tests/rfc3977/body_workflow.rs +++ b/tests/rfc3977/body_workflow.rs @@ -681,6 +681,186 @@ async fn expect_pass_through_borrows_initial_backend_bytes_with_auth( Ok(()) } +#[tokio::test] +async fn test_per_command_body_window_sends_two_requests_before_first_reply() -> Result<()> { + let backend_listener = TcpListener::bind("127.0.0.1:0").await?; + let backend_port = backend_listener.local_addr()?.port(); + let (connection_count, seen_commands) = spawn_delayed_pipeline_backend( + backend_listener, + &[ + ( + "BODY ", + "222 1 \r\nwindow-1-line\r\n.\r\n", + ), + ( + "BODY ", + "222 2 \r\nwindow-2-line\r\n.\r\n", + ), + ], + &[], + ); + + let proxy_port = spawn_proxy_with_config( + pipeline_backend_config(backend_port, "BodyUpstreamWindow"), + RoutingMode::PerCommand, + ) + .await?; + let client = connect_and_read_greeting(proxy_port).await?; + let (read_half, mut write_half) = client.into_split(); + let mut reader = BufReader::new(read_half); + + write_commands( + &mut write_half, + &["BODY \r\n", "BODY \r\n"], + ) + .await?; + + let responses = read_multiline_responses(&mut reader, 2).await?; + assert_eq!( + responses, + vec![ + ( + "222 1 \r\n".to_string(), + vec!["window-1-line\r\n".to_string()], + ), + ( + "222 2 \r\n".to_string(), + vec!["window-2-line\r\n".to_string()], + ), + ] + ); + assert_eq!(connection_count.load(Ordering::SeqCst), 1); + assert_eq!( + seen_commands + .lock() + .await + .iter() + .filter(|command| command.starts_with("BODY ")) + .count(), + 2, + ); + + Ok(()) +} + +#[tokio::test] +async fn test_per_command_body_window_keeps_430_before_following_success() -> Result<()> { + let backend_listener = TcpListener::bind("127.0.0.1:0").await?; + let backend_port = backend_listener.local_addr()?.port(); + let (_connection_count, _seen_commands) = spawn_delayed_pipeline_backend( + backend_listener, + &[ + ("BODY ", "430 No such article\r\n"), + ( + "BODY ", + "222 2 \r\nwindow-present-line\r\n.\r\n", + ), + ], + &[], + ); + + let proxy_port = spawn_proxy_with_config( + pipeline_backend_config(backend_port, "BodyUpstreamWindow430"), + RoutingMode::PerCommand, + ) + .await?; + let client = connect_and_read_greeting(proxy_port).await?; + let (read_half, mut write_half) = client.into_split(); + let mut reader = BufReader::new(read_half); + + write_commands( + &mut write_half, + &[ + "BODY \r\n", + "BODY \r\n", + ], + ) + .await?; + + assert_eq!( + read_line(&mut reader, "missing presend response").await?, + "430 No such article\r\n", + ); + assert_eq!( + read_multiline_response(&mut reader).await?, + ( + "222 2 \r\n".to_string(), + vec!["window-present-line\r\n".to_string()], + ), + ); + + Ok(()) +} + +#[tokio::test] +async fn test_per_command_body_window_keeps_cached_430_after_earlier_success() -> Result<()> { + let backend_listener = TcpListener::bind("127.0.0.1:0").await?; + let backend_port = backend_listener.local_addr()?.port(); + let (_connection_count, seen_commands) = spawn_delayed_pipeline_backend( + backend_listener, + &[( + "BODY ", + "222 2 \r\nwindow-present-line\r\n.\r\n", + )], + &[( + "BODY ", + "430 No such article\r\n", + )], + ); + + let proxy_port = spawn_proxy_with_config( + pipeline_backend_config(backend_port, "BodyUpstreamWindowCached430"), + RoutingMode::PerCommand, + ) + .await?; + let client = connect_and_read_greeting(proxy_port).await?; + let (read_half, mut write_half) = client.into_split(); + let mut reader = BufReader::new(read_half); + + write_commands( + &mut write_half, + &["BODY \r\n"], + ) + .await?; + assert_eq!( + read_line(&mut reader, "cache warmup 430").await?, + "430 No such article\r\n", + ); + + write_commands( + &mut write_half, + &[ + "BODY \r\n", + "BODY \r\n", + ], + ) + .await?; + + assert_eq!( + read_multiline_response(&mut reader).await?, + ( + "222 2 \r\n".to_string(), + vec!["window-present-line\r\n".to_string()], + ), + ); + assert_eq!( + read_line(&mut reader, "cached missing response").await?, + "430 No such article\r\n", + ); + + let seen_commands = seen_commands.lock().await; + assert_eq!( + seen_commands + .iter() + .filter(|command| command.as_str() == "BODY ") + .count(), + 1, + "known-missing request must not be sent upstream again", + ); + + Ok(()) +} + #[tokio::test] async fn test_body_pipelining_pairs_four_responses_on_single_backend_connection() -> Result<()> { let backend_listener = TcpListener::bind("127.0.0.1:0").await?;