From a3744443ea10d0fcdfc59212a971a4361ed76434 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Sat, 5 Sep 2026 14:47:02 -0600 Subject: [PATCH 1/3] protocol: reject transport-changing extensions --- src/command/handler.rs | 28 ++++++++++++++++++++++++++++ src/protocol/request.rs | 7 ++++++- 2 files changed, 34 insertions(+), 1 deletion(-) diff --git a/src/command/handler.rs b/src/command/handler.rs index b83bee6b..c115ec64 100644 --- a/src/command/handler.rs +++ b/src/command/handler.rs @@ -246,6 +246,10 @@ const STATEFUL_REJECT: RejectResponse = RejectResponse::new( StatusCode::new(codes::FEATURE_NOT_SUPPORTED), "503 Feature not supported in stateless proxy mode\r\n", ); +const TRANSPORT_REJECT: RejectResponse = RejectResponse::new( + StatusCode::new(codes::FEATURE_NOT_SUPPORTED), + "503 Transport-changing command not supported\r\n", +); const fn wire_status(wire: &str) -> u16 { let bytes = wire.as_bytes(); @@ -368,6 +372,7 @@ fn rejection_for(request: &RequestContext) -> RejectResponse { match request.kind() { RequestKind::Post => POST_REJECT, RequestKind::Ihave => TRANSIT_REJECT, + RequestKind::Compress => TRANSPORT_REJECT, _ => match request.route_class() { RequestRouteClass::Stateful => STATEFUL_REJECT, _ => TRANSIT_REJECT, @@ -564,6 +569,29 @@ mod tests { ); } + #[test] + fn compress_is_rejected_without_forwarding_in_any_routing_mode() { + let request = RequestContext::parse(b"COMPRESS DEFLATE\r\n").expect("valid command"); + assert_eq!(request.kind(), RequestKind::Compress); + assert_eq!(request.route_class(), RequestRouteClass::Reject); + + for routing_mode in [ + crate::config::RoutingMode::PerCommand, + crate::config::RoutingMode::Hybrid, + crate::config::RoutingMode::Stateful, + ] { + let plan = CommandHandler::classify_request( + &request, + AuthenticationAccess::Authenticated, + routing_mode, + ); + assert!( + matches!(plan, CommandPlan::Reject(response) if response.status().as_u16() == 503), + "COMPRESS must be rejected in {routing_mode:?}: {plan:?}" + ); + } + } + /// Bug 2 regression test: RFC 4643 §2.3.1 — AUTHINFO is case-insensitive. /// /// Before fix: the username/password extractor only stripped exact "AUTHINFO USER" or diff --git a/src/protocol/request.rs b/src/protocol/request.rs index 3887f04c..e35cb9df 100644 --- a/src/protocol/request.rs +++ b/src/protocol/request.rs @@ -39,6 +39,7 @@ pub enum RequestKind { TakeThis, AuthInfo, StartTls, + Compress, Unknown, } @@ -834,7 +835,8 @@ const fn route_class(kind: RequestKind, has_message_id: bool) -> RequestRouteCla | RequestKind::Ihave | RequestKind::Check | RequestKind::TakeThis - | RequestKind::StartTls => RequestRouteClass::Reject, + | RequestKind::StartTls + | RequestKind::Compress => RequestRouteClass::Reject, RequestKind::Article | RequestKind::Body | RequestKind::Head | RequestKind::Stat if has_message_id => { @@ -917,6 +919,7 @@ const fn classify_verb(verb: &[u8]) -> RequestKind { }, 8 => { b"AUTHINFO" => RequestKind::AuthInfo, + b"COMPRESS" => RequestKind::Compress, b"STARTTLS" => RequestKind::StartTls, b"TAKETHIS" => RequestKind::TakeThis, }, @@ -1412,6 +1415,7 @@ mod tests { ("TAKETHIS \r\n", RequestKind::TakeThis), ("AUTHINFO USER test\r\n", RequestKind::AuthInfo), ("STARTTLS\r\n", RequestKind::StartTls), + ("COMPRESS DEFLATE\r\n", RequestKind::Compress), ]; for (line, expected) in cases { @@ -1438,6 +1442,7 @@ mod tests { ("CHECK \r\n", RequestRouteClass::Reject), ("TAKETHIS \r\n", RequestRouteClass::Reject), ("STARTTLS\r\n", RequestRouteClass::Reject), + ("COMPRESS DEFLATE\r\n", RequestRouteClass::Reject), ("XFOO arg\r\n", RequestRouteClass::Stateful), ]; From 760b71180f27b1f4526655c4bffa7df54c611676 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Sat, 5 Sep 2026 17:14:15 -0600 Subject: [PATCH 2/3] fix: classify STARTTLS as transport rejection --- src/command/handler.rs | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/src/command/handler.rs b/src/command/handler.rs index c115ec64..50c05bc6 100644 --- a/src/command/handler.rs +++ b/src/command/handler.rs @@ -372,7 +372,7 @@ fn rejection_for(request: &RequestContext) -> RejectResponse { match request.kind() { RequestKind::Post => POST_REJECT, RequestKind::Ihave => TRANSIT_REJECT, - RequestKind::Compress => TRANSPORT_REJECT, + RequestKind::Compress | RequestKind::StartTls => TRANSPORT_REJECT, _ => match request.route_class() { RequestRouteClass::Stateful => STATEFUL_REJECT, _ => TRANSIT_REJECT, @@ -585,9 +585,13 @@ mod tests { AuthenticationAccess::Authenticated, routing_mode, ); - assert!( - matches!(plan, CommandPlan::Reject(response) if response.status().as_u16() == 503), - "COMPRESS must be rejected in {routing_mode:?}: {plan:?}" + let CommandPlan::Reject(response) = plan else { + panic!("COMPRESS must be rejected in {routing_mode:?}: {plan:?}"); + }; + assert_eq!(response.status().as_u16(), 503); + assert_eq!( + response.to_string(), + "503 Transport-changing command not supported\r\n" ); } } From 6a3df59bccbb0c41a7bc16847bccb98bf56ec800 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Fri, 11 Sep 2026 12:13:20 -0600 Subject: [PATCH 3/3] Benchmark metadata and payload cache updates --- Cargo.toml | 7 ++ benches/cache_metadata_payload.rs | 120 ++++++++++++++++++++++++++++++ 2 files changed, 127 insertions(+) create mode 100644 benches/cache_metadata_payload.rs diff --git a/Cargo.toml b/Cargo.toml index c5c378f8..a499cda6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -176,17 +176,24 @@ harness = false name = "cache_ingest" harness = false +[[bench]] +name = "cache_metadata_payload" +harness = false + [[bench]] name = "end_to_end_proxy" harness = false [profile.dev] +debug = 0 # Use the profiling/bench profiles when symbols are needed +incremental = false # Avoid retaining a second full set of object files split-debuginfo = "unpacked" # Faster linking — don't bundle debuginfo into binary [profile.dev.package."*"] opt-level = 1 # Compile dependencies with optimizations in dev mode # Huge runtime speedup for rustls/ring/foyer/moka # Minimal compile-time cost (deps cached after first build) +debug = 0 # Keep dependency artifacts compact in the dev profile [profile.release] debug = false # No debug symbols for smaller binaries diff --git a/benches/cache_metadata_payload.rs b/benches/cache_metadata_payload.rs new file mode 100644 index 00000000..22826ac5 --- /dev/null +++ b/benches/cache_metadata_payload.rs @@ -0,0 +1,120 @@ +//! Mixed hybrid-cache workload for metadata and retained-payload updates. +//! +//! Run with: `cargo bench --bench cache_metadata_payload` + +use divan::Bencher; +use nntp_proxy::cache::{HybridCacheConfig, UnifiedCache}; +use nntp_proxy::protocol::StatusCode; +use nntp_proxy::types::{BackendId, MessageId}; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; +use tempfile::{TempDir, tempdir}; + +const ARTICLE_BODY: &str = "x"; + +fn main() { + divan::main(); +} + +fn benchmark_cache() -> (tokio::runtime::Runtime, TempDir, UnifiedCache) { + let runtime = tokio::runtime::Runtime::new().expect("benchmark runtime"); + let directory = tempdir().expect("benchmark cache directory"); + let config = HybridCacheConfig { + memory_capacity: 4 * 1024 * 1024, + disk_capacity: 64 * 1024 * 1024, + disk_path: directory.path().to_path_buf(), + ttl: Duration::from_secs(300), + compression: nntp_proxy::config::CompressionCodec::None, + shards: 16, + }; + let cache = runtime + .block_on(UnifiedCache::hybrid(config)) + .expect("hybrid cache"); + (runtime, directory, cache) +} + +fn message_id(sequence: u64) -> MessageId<'static> { + MessageId::new(format!("")).expect("benchmark message ID") +} + +fn article_response(sequence: u64) -> Vec { + format!( + "220 42 \r\nSubject: Benchmark\r\n\r\n{ARTICLE_BODY}\r\n.\r\n" + ) + .into_bytes() +} + +#[divan::bench(sample_count = 20, sample_size = 1)] +fn metadata_only_updates(bencher: Bencher) { + let (runtime, _directory, cache) = benchmark_cache(); + let sequence = AtomicU64::new(0); + + bencher + .with_inputs(|| message_id(sequence.fetch_add(1, Ordering::Relaxed))) + .bench_values(|id| { + runtime.block_on(cache.record_backend_has_status( + id, + StatusCode::new(223), + BackendId::from_index(0), + 0.into(), + )); + }); + + runtime + .block_on(cache.close()) + .expect("close benchmark cache"); +} + +#[divan::bench(sample_count = 20, sample_size = 1)] +fn retained_payload_updates(bencher: Bencher) { + let (runtime, _directory, cache) = benchmark_cache(); + let sequence = AtomicU64::new(0); + + bencher + .with_inputs(|| { + let sequence = sequence.fetch_add(1, Ordering::Relaxed); + (message_id(sequence), article_response(sequence)) + }) + .bench_values(|(id, response)| { + runtime.block_on(cache.upsert_ingest(id, response, BackendId::from_index(0), 0.into())); + }); + + runtime + .block_on(cache.close()) + .expect("close benchmark cache"); +} + +#[divan::bench(sample_count = 20, sample_size = 1)] +fn mixed_metadata_and_payload_updates(bencher: Bencher) { + let (runtime, _directory, cache) = benchmark_cache(); + let sequence = AtomicU64::new(0); + + bencher + .with_inputs(|| { + let sequence = sequence.fetch_add(2, Ordering::Relaxed); + ( + message_id(sequence), + message_id(sequence + 1), + article_response(sequence + 1), + ) + }) + .bench_values(|(metadata_id, payload_id, response)| { + runtime.block_on(async { + cache + .record_backend_has_status( + metadata_id, + StatusCode::new(223), + BackendId::from_index(0), + 0.into(), + ) + .await; + cache + .upsert_ingest(payload_id, response, BackendId::from_index(0), 0.into()) + .await; + }); + }); + + runtime + .block_on(cache.close()) + .expect("close benchmark cache"); +}