diff --git a/Cargo.lock b/Cargo.lock index 02c1fa1..d81d679 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1,6 +1,6 @@ # This file is automatically @generated by Cargo. # It is not intended for manual editing. -version = 3 +version = 4 [[package]] name = "aho-corasick" @@ -89,6 +89,12 @@ dependencies = [ "tracing", ] +[[package]] +name = "bitflags" +version = "2.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b4388bee8683e3d04af747c73422af53102d2bd24d9eadb6cbc100baef4b43f8" + [[package]] name = "bytes" version = "1.11.1" @@ -118,6 +124,7 @@ version = "0.1.0" dependencies = [ "axum", "codag-drain", + "dashmap", "http-body-util", "serde", "serde_json", @@ -148,6 +155,19 @@ dependencies = [ "memchr", ] +[[package]] +name = "dashmap" +version = "5.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "978747c1d849a7d2ee5e8adc0159961c48fb7e5db2f06af6723b80123bb53856" +dependencies = [ + "cfg-if", + "hashbrown 0.14.5", + "lock_api", + "once_cell", + "parking_lot_core", +] + [[package]] name = "drain3_rust" version = "0.1.0" @@ -211,6 +231,12 @@ dependencies = [ "slab", ] +[[package]] +name = "hashbrown" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" + [[package]] name = "hashbrown" version = "0.15.5" @@ -320,6 +346,15 @@ version = "0.2.186" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "68ab91017fe16c622486840e4c83c9a37afeff978bd239b5293d61ece587de66" +[[package]] +name = "lock_api" +version = "0.4.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "224399e74b87b5f3557511d98dff8b14089b3dadafcab6bb93eab67d3aace965" +dependencies = [ + "scopeguard", +] + [[package]] name = "log" version = "0.4.29" @@ -332,7 +367,7 @@ version = "0.12.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "234cf4f4a04dc1f57e24b96cc0cd600cf2af460d4161ac5ecdd0af8e1f3b2a38" dependencies = [ - "hashbrown", + "hashbrown 0.15.5", ] [[package]] @@ -388,6 +423,19 @@ version = "1.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" +[[package]] +name = "parking_lot_core" +version = "0.9.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2621685985a2ebf1c516881c026032ac7deafcda1a2c9b7850dc81e3dfcb64c1" +dependencies = [ + "cfg-if", + "libc", + "redox_syscall", + "smallvec", + "windows-link", +] + [[package]] name = "percent-encoding" version = "2.3.2" @@ -418,6 +466,15 @@ dependencies = [ "proc-macro2", ] +[[package]] +name = "redox_syscall" +version = "0.5.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed2bf2547551a7053d6fdfafda3f938979645c44812fbfcda098faae3f1a362d" +dependencies = [ + "bitflags", +] + [[package]] name = "regex" version = "1.12.3" @@ -459,6 +516,12 @@ version = "1.0.23" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" +[[package]] +name = "scopeguard" +version = "1.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" + [[package]] name = "serde" version = "1.0.228" diff --git a/examples/server/Cargo.toml b/examples/server/Cargo.toml index e4400c1..725fad0 100644 --- a/examples/server/Cargo.toml +++ b/examples/server/Cargo.toml @@ -14,6 +14,7 @@ serde = { version = "1", features = ["derive"] } serde_json = "1" tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter"] } +dashmap = "5.5" [dev-dependencies] tower = { version = "0.5", features = ["util"] } diff --git a/examples/server/src/main.rs b/examples/server/src/main.rs index 29fcc10..b67016d 100644 --- a/examples/server/src/main.rs +++ b/examples/server/src/main.rs @@ -22,12 +22,15 @@ async fn ttl_sweeper(state: AppState) { loop { iv.tick().await; let now = Instant::now(); - let mut sessions = state.sessions.write().await; - let before = sessions.len(); - sessions.retain(|_, e| now.duration_since(e.last_touch) < SESSION_TTL); - let evicted = before - sessions.len(); + let before = state.sessions.len(); + state.sessions.retain(|_, e| { + let last_touch = *e.last_touch.lock().unwrap_or_else(|e| e.into_inner()); + now.duration_since(last_touch) < SESSION_TTL + }); + let after = state.sessions.len(); + let evicted = before.saturating_sub(after); if evicted > 0 { - tracing::info!(evicted, remaining = sessions.len(), "ttl sweep"); + tracing::info!(evicted, remaining = after, "ttl sweep"); } } } diff --git a/examples/server/src/routes.rs b/examples/server/src/routes.rs index 17f1686..61b130d 100644 --- a/examples/server/src/routes.rs +++ b/examples/server/src/routes.rs @@ -223,13 +223,24 @@ async fn ingest( }; let ingested = new_lines.len(); - let mut sessions = state.sessions.write().await; - let entry = sessions.entry(id.clone()).or_insert_with(SessionEntry::new); + // Fast path: DashMap::get only takes a read lock on the shard, maximizing concurrency + // for existing sessions. + let entry = if let Some(e) = state.sessions.get(&id) { + e.clone() + } else { + // Slow path: DashMap::entry takes a write lock on the shard to insert. + state.sessions.entry(id.clone()).or_insert_with(|| std::sync::Arc::new(SessionEntry::new())).value().clone() + }; + + let mut index = entry.index.write().await; for l in new_lines { - entry.index.push(l); + index.push(l); + } + { + let mut lock = entry.last_touch.lock().unwrap_or_else(|e| e.into_inner()); + *lock = Instant::now(); } - entry.last_touch = Instant::now(); - let total = entry.index.len(); + let total = index.len(); tracing::info!(session = %id, ingested, total, "ingest"); Json(json!({ "ingested": ingested, "total": total })).into_response() @@ -244,11 +255,12 @@ async fn templates( Ok(cfg) => cfg, Err(resp) => return resp, }; - let sessions = state.sessions.read().await; - let Some(entry) = sessions.get(&id) else { - return (StatusCode::NOT_FOUND, format!("no such session: {id}")).into_response(); + let entry = match state.sessions.get(&id) { + Some(e) => e.clone(), + None => return (StatusCode::NOT_FOUND, format!("no such session: {id}")).into_response(), }; - let result = entry.index.templates_with(&cfg); + let index = entry.index.read().await; + let result = index.templates_with(&cfg); if wants_json(q.format.as_deref()) { Json(result).into_response() } else { diff --git a/examples/server/src/session.rs b/examples/server/src/session.rs index 0c1352d..c8e062f 100644 --- a/examples/server/src/session.rs +++ b/examples/server/src/session.rs @@ -5,12 +5,12 @@ //! lives entirely in memory; sessions idle longer than [`SESSION_TTL`] are //! evicted by a background task (see `main.rs`). -use std::collections::HashMap; use std::env; use std::sync::Arc; use std::time::{Duration, Instant}; use codag_drain::{TemplateIndex, TemplaterConfig}; +use dashmap::DashMap; use tokio::sync::{RwLock, Semaphore}; pub use codag_drain::{parse_body, parse_json_line, parse_line, BodyFormat}; @@ -67,8 +67,8 @@ fn env_u64(name: &str, default: u64) -> u64 { /// One live session: a streaming index plus its last-touch timestamp. #[derive(Debug)] pub struct SessionEntry { - pub index: TemplateIndex, - pub last_touch: Instant, + pub index: RwLock, + pub last_touch: std::sync::Mutex, } impl SessionEntry { @@ -80,8 +80,8 @@ impl SessionEntry { impl Default for SessionEntry { fn default() -> Self { SessionEntry { - index: TemplateIndex::new(TemplaterConfig::default()), - last_touch: Instant::now(), + index: RwLock::new(TemplateIndex::new(TemplaterConfig::default())), + last_touch: std::sync::Mutex::new(Instant::now()), } } } @@ -89,7 +89,7 @@ impl Default for SessionEntry { /// Shared, cloneable application state. #[derive(Clone)] pub struct AppState { - pub sessions: Arc>>, + pub sessions: Arc>>, pub limits: Arc, pub template_slots: Arc, pub auth_token: Option>, @@ -100,7 +100,7 @@ impl AppState { let limits = Limits::from_env(); let max_inflight = limits.max_inflight; AppState { - sessions: Arc::new(RwLock::new(HashMap::new())), + sessions: Arc::new(DashMap::new()), limits: Arc::new(limits), template_slots: Arc::new(Semaphore::new(max_inflight)), auth_token: auth_token_from_env(), diff --git a/examples/server/tests/integration.rs b/examples/server/tests/integration.rs index 5f954bc..5468fc9 100644 --- a/examples/server/tests/integration.rs +++ b/examples/server/tests/integration.rs @@ -4,11 +4,11 @@ use axum::body::Body; use axum::http::{HeaderValue, Request, StatusCode}; use http_body_util::BodyExt; -use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; -use tokio::sync::{RwLock, Semaphore}; +use tokio::sync::Semaphore; use tower::ServiceExt; +use dashmap::DashMap; use codag_drain::{template_logs, LogLine, TemplaterConfig}; use codag_drain_server::routes; @@ -48,7 +48,7 @@ fn limited_app(limits: Limits) -> axum::Router { fn state_with_limits(limits: Limits) -> AppState { let max_inflight = limits.max_inflight; AppState { - sessions: Arc::new(RwLock::new(HashMap::new())), + sessions: Arc::new(DashMap::new()), limits: Arc::new(limits), template_slots: Arc::new(Semaphore::new(max_inflight)), auth_token: Some(Arc::::from(TEST_TOKEN)), @@ -105,7 +105,7 @@ async fn v1_routes_fail_closed_without_configured_auth_token() { let limits = Limits::from_env(); let max_inflight = limits.max_inflight; let state = AppState { - sessions: Arc::new(RwLock::new(HashMap::new())), + sessions: Arc::new(DashMap::new()), limits: Arc::new(limits), template_slots: Arc::new(Semaphore::new(max_inflight)), auth_token: None, @@ -345,7 +345,7 @@ async fn concurrency_overflow_returns_429() { template_timeout: Duration::from_secs(1), }; let state = AppState { - sessions: Arc::new(RwLock::new(HashMap::new())), + sessions: Arc::new(DashMap::new()), limits: Arc::new(limits), template_slots: Arc::new(Semaphore::new(1)), auth_token: Some(Arc::::from(TEST_TOKEN)),