This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-pack / src / cache.rs
16 kB 511 lines
1use std::collections::HashMap; 2use std::path::{Path, PathBuf}; 3use std::sync::{Arc, Mutex}; 4use std::time::Duration; 5 6use axum::body::Bytes; 7use knot_cache::{Cache, EntryCount, Lru, Reclaimable, Weight}; 8use knot_runtime::Clock; 9use tokio::sync::watch; 10 11const MAX_ENTRIES: usize = 1024; 12 13knot_types::scalar_newtype! { 14 pub struct MaxEntryBytes(usize); 15 pub struct MaxCacheBytes(usize); 16} 17 18#[derive(Debug, Clone, Copy)] 19pub struct CacheConfig { 20 pub enabled: bool, 21 pub ttl: Duration, 22 pub max_entry_bytes: MaxEntryBytes, 23 pub max_total_bytes: MaxCacheBytes, 24} 25 26impl Default for CacheConfig { 27 fn default() -> Self { 28 Self { 29 enabled: true, 30 ttl: Duration::from_secs(60), 31 max_entry_bytes: MaxEntryBytes::new(64 * 1024 * 1024), 32 max_total_bytes: MaxCacheBytes::new(512 * 1024 * 1024), 33 } 34 } 35} 36 37#[derive(Clone, PartialEq, Eq, Hash)] 38pub(crate) struct RequestKey { 39 objects_dir: PathBuf, 40 digest: [u8; 32], 41} 42 43impl RequestKey { 44 pub(crate) fn new(objects_dir: &Path, ref_token: &crate::ids::RefsDigest, body: &[u8]) -> Self { 45 let mut hasher = gix::hash::hasher(gix::hash::Kind::Sha256); 46 hasher.update(ref_token.as_bytes()); 47 hasher.update(body); 48 let id = hasher.try_finalize().expect("sha256 digest finalizes"); 49 let mut digest = [0u8; 32]; 50 digest.copy_from_slice(id.as_slice()); 51 Self { 52 objects_dir: objects_dir.to_path_buf(), 53 digest, 54 } 55 } 56} 57 58#[derive(Clone)] 59pub(crate) enum Signal { 60 Pending, 61 Ready(Bytes), 62 Retry, 63 Regenerate, 64} 65 66#[derive(Clone)] 67enum Slot { 68 Ready(Bytes), 69 TooLarge, 70} 71 72fn weigh(slot: &Slot) -> Weight { 73 match slot { 74 Slot::Ready(bytes) => Weight::new(bytes.len() as u64), 75 Slot::TooLarge => Weight::new(0), 76 } 77} 78 79enum Settlement { 80 Ready(Bytes), 81 TooLarge, 82 Regenerate, 83 Retry, 84} 85 86pub(crate) struct PackCache { 87 config: CacheConfig, 88 store: Lru<RequestKey, Slot, Arc<dyn Clock>>, 89 inflight: Mutex<HashMap<RequestKey, watch::Sender<Signal>>>, 90} 91 92pub(crate) enum Decision { 93 Serve(Bytes), 94 Await(watch::Receiver<Signal>), 95 Lead(Lease), 96 Stream, 97 Off, 98} 99 100impl PackCache { 101 pub(crate) fn new(mut config: CacheConfig, clock: Arc<dyn Clock>) -> Arc<Self> { 102 config.max_entry_bytes = MaxEntryBytes::new( 103 config 104 .max_entry_bytes 105 .get() 106 .min(config.max_total_bytes.get()), 107 ); 108 let store = Lru::by_weight_with_ttl( 109 Weight::new(config.max_total_bytes.get() as u64), 110 config.ttl, 111 clock, 112 weigh, 113 ) 114 .with_entry_cap(EntryCount::new(MAX_ENTRIES as u64)); 115 let cache = Arc::new(Self { 116 config, 117 store, 118 inflight: Mutex::new(HashMap::new()), 119 }); 120 knot_cache::register(&cache); 121 cache 122 } 123 124 pub(crate) fn decide(self: &Arc<Self>, key: RequestKey) -> Decision { 125 if !self.config.enabled { 126 return Decision::Off; 127 } 128 if let Some(decision) = self.serve_cached(&key) { 129 return decision; 130 } 131 let mut inflight = self.inflight_lock(); 132 if let Some(sender) = inflight.get(&key) { 133 return Decision::Await(sender.subscribe()); 134 } 135 // `store` & `inflight` are separate locks, 136 // such that a leader can finish settling in the gap between 137 // a fast-path miss and this lock acquisition, 138 // by which point it has inserted the pack and taken its sender away. 139 // Reading the store a second time will cover that gap, 140 // since `settle` always inserts before it removes. 141 // 142 // Without it, 143 // the race-losing caller would elect itself and rebuild a pack 144 // the store does in fact already have. 145 if let Some(decision) = self.serve_cached(&key) { 146 return decision; 147 } 148 let (sender, _) = watch::channel(Signal::Pending); 149 inflight.insert(key.clone(), sender); 150 Decision::Lead(Lease { 151 key, 152 cache: Arc::clone(self), 153 max_entry_bytes: self.config.max_entry_bytes, 154 settled: false, 155 }) 156 } 157 158 fn serve_cached(&self, key: &RequestKey) -> Option<Decision> { 159 match self.store.get(key) { 160 Some(Slot::Ready(bytes)) => Some(Decision::Serve(bytes)), 161 Some(Slot::TooLarge) => Some(Decision::Stream), 162 None => None, 163 } 164 } 165 166 fn settle(&self, key: &RequestKey, settlement: Settlement) { 167 let signal = match settlement { 168 Settlement::Ready(bytes) => { 169 self.store.insert(key.clone(), Slot::Ready(bytes.clone())); 170 Signal::Ready(bytes) 171 } 172 Settlement::TooLarge => { 173 self.store.insert(key.clone(), Slot::TooLarge); 174 Signal::Regenerate 175 } 176 Settlement::Regenerate => Signal::Regenerate, 177 Settlement::Retry => Signal::Retry, 178 }; 179 if let Some(sender) = self.inflight_lock().remove(key) { 180 let _ = sender.send(signal); 181 } 182 } 183 184 fn inflight_lock( 185 &self, 186 ) -> std::sync::MutexGuard<'_, HashMap<RequestKey, watch::Sender<Signal>>> { 187 self.inflight 188 .lock() 189 .unwrap_or_else(std::sync::PoisonError::into_inner) 190 } 191 192 #[cfg(test)] 193 fn entry_count(&self) -> u64 { 194 self.store.entry_count().get() 195 } 196} 197 198impl Reclaimable for PackCache { 199 fn footprint(&self) -> Weight { 200 self.store.footprint() 201 } 202 203 fn reclaim(&self) { 204 self.store.reclaim(); 205 } 206} 207 208pub(crate) struct Lease { 209 key: RequestKey, 210 cache: Arc<PackCache>, 211 max_entry_bytes: MaxEntryBytes, 212 settled: bool, 213} 214 215impl Lease { 216 pub(crate) fn max_entry_bytes(&self) -> MaxEntryBytes { 217 self.max_entry_bytes 218 } 219 220 pub(crate) fn ready(mut self, bytes: Bytes) { 221 self.cache.settle(&self.key, Settlement::Ready(bytes)); 222 self.settled = true; 223 } 224 225 pub(crate) fn too_large(mut self) { 226 self.cache.settle(&self.key, Settlement::TooLarge); 227 self.settled = true; 228 } 229 230 pub(crate) fn regenerate(mut self) { 231 self.cache.settle(&self.key, Settlement::Regenerate); 232 self.settled = true; 233 } 234 235 pub(crate) fn retry(mut self) { 236 self.cache.settle(&self.key, Settlement::Retry); 237 self.settled = true; 238 } 239} 240 241impl Drop for Lease { 242 fn drop(&mut self) { 243 if !self.settled { 244 self.cache.settle(&self.key, Settlement::Retry); 245 } 246 } 247} 248 249pub(crate) enum Resolved { 250 Bytes(Bytes), 251 Retry, 252 Regenerate, 253} 254 255pub(crate) async fn wait(mut receiver: watch::Receiver<Signal>) -> Resolved { 256 let resolved = classify(&receiver.borrow_and_update()); 257 match resolved { 258 Some(resolved) => resolved, 259 None => match receiver.changed().await { 260 Err(_) => Resolved::Regenerate, 261 Ok(()) => Box::pin(wait(receiver)).await, 262 }, 263 } 264} 265 266fn classify(signal: &Signal) -> Option<Resolved> { 267 match signal { 268 Signal::Pending => None, 269 Signal::Ready(bytes) => Some(Resolved::Bytes(bytes.clone())), 270 Signal::Retry => Some(Resolved::Retry), 271 Signal::Regenerate => Some(Resolved::Regenerate), 272 } 273} 274 275pub(crate) enum Capture { 276 Buffering { buffer: Vec<u8>, limit: usize }, 277 Overflow, 278 Off, 279} 280 281impl Capture { 282 pub(crate) fn new(limit: Option<MaxEntryBytes>) -> Self { 283 match limit { 284 Some(limit) => Capture::Buffering { 285 buffer: Vec::new(), 286 limit: limit.get(), 287 }, 288 None => Capture::Off, 289 } 290 } 291 292 pub(crate) fn record(&mut self, chunk: &[u8]) { 293 match self { 294 Capture::Buffering { buffer, limit } if buffer.len() + chunk.len() <= *limit => { 295 buffer.extend_from_slice(chunk) 296 } 297 Capture::Buffering { .. } => *self = Capture::Overflow, 298 _ => {} 299 } 300 } 301 302 pub(crate) fn into_bytes(self) -> Option<Vec<u8>> { 303 match self { 304 Capture::Buffering { buffer, .. } => Some(buffer), 305 _ => None, 306 } 307 } 308} 309 310#[cfg(test)] 311mod tests { 312 use knot_runtime::{ManualClock, SystemClock, UnixMicros}; 313 314 use super::*; 315 316 fn built(config: CacheConfig) -> Arc<PackCache> { 317 PackCache::new(config, Arc::new(SystemClock)) 318 } 319 320 fn key(body: &[u8]) -> RequestKey { 321 let mut token = [0u8; 32]; 322 token[..4].copy_from_slice(b"refs"); 323 RequestKey::new( 324 Path::new("/scan/did:plc:squid/objects"), 325 &crate::ids::RefsDigest::new(token), 326 body, 327 ) 328 } 329 330 fn lead(cache: &Arc<PackCache>, body: &[u8]) -> Lease { 331 match cache.decide(key(body)) { 332 Decision::Lead(lease) => lease, 333 _ => panic!("the first request for a key leads"), 334 } 335 } 336 337 #[tokio::test] 338 async fn a_second_identical_request_serves_the_cached_pack() { 339 let cache = built(CacheConfig::default()); 340 lead(&cache, b"want").ready(Bytes::from_static(b"PACK-bytes")); 341 match cache.decide(key(b"want")) { 342 Decision::Serve(bytes) => assert_eq!(bytes.as_ref(), b"PACK-bytes"), 343 _ => panic!("second identical request hits the cache"), 344 } 345 } 346 347 #[tokio::test] 348 async fn a_concurrent_request_awaits_the_leader_then_shares_its_bytes() { 349 let cache = built(CacheConfig::default()); 350 let leader = lead(&cache, b"clone"); 351 let receiver = match cache.decide(key(b"clone")) { 352 Decision::Await(receiver) => receiver, 353 _ => panic!("concurrent request awaits the inflight leader"), 354 }; 355 leader.ready(Bytes::from_static(b"shared")); 356 match wait(receiver).await { 357 Resolved::Bytes(bytes) => assert_eq!(bytes.as_ref(), b"shared"), 358 _ => panic!("the follower shares the leader's bytes"), 359 } 360 } 361 362 #[tokio::test] 363 async fn a_regenerating_leader_wakes_followers_to_regenerate() { 364 let cache = built(CacheConfig::default()); 365 let leader = lead(&cache, b"err"); 366 let receiver = match cache.decide(key(b"err")) { 367 Decision::Await(receiver) => receiver, 368 _ => panic!("follower awaits"), 369 }; 370 leader.regenerate(); 371 assert!( 372 matches!(wait(receiver).await, Resolved::Regenerate), 373 "an errored leader tells the follower to regenerate in parallel" 374 ); 375 match cache.decide(key(b"err")) { 376 Decision::Lead(_) => {} 377 _ => panic!("a regenerated key isn't remembered, so the next request leads afresh"), 378 } 379 } 380 381 #[tokio::test] 382 async fn a_retrying_leader_tells_followers_to_re_elect() { 383 let cache = built(CacheConfig::default()); 384 let leader = lead(&cache, b"vanish"); 385 let receiver = match cache.decide(key(b"vanish")) { 386 Decision::Await(receiver) => receiver, 387 _ => panic!("follower awaits"), 388 }; 389 leader.retry(); 390 assert!( 391 matches!(wait(receiver).await, Resolved::Retry), 392 "a vanished leader tells the follower to re-elect a fresh leader" 393 ); 394 match cache.decide(key(b"vanish")) { 395 Decision::Lead(_) => {} 396 _ => panic!("a retried key isn't remembered, so the next request leads afresh"), 397 } 398 } 399 400 #[tokio::test] 401 async fn a_dropped_lease_re_elects_rather_than_stranding_the_follower() { 402 let cache = built(CacheConfig::default()); 403 let leader = lead(&cache, b"dropped"); 404 let receiver = match cache.decide(key(b"dropped")) { 405 Decision::Await(receiver) => receiver, 406 _ => panic!("follower awaits"), 407 }; 408 drop(leader); 409 assert!( 410 matches!(wait(receiver).await, Resolved::Retry), 411 "an unsettled lease that drops re-elects instead of stranding the follower" 412 ); 413 } 414 415 #[tokio::test] 416 async fn an_oversized_leader_marks_the_key_for_direct_streaming() { 417 let cache = built(CacheConfig::default()); 418 let leader = lead(&cache, b"huge"); 419 let receiver = match cache.decide(key(b"huge")) { 420 Decision::Await(receiver) => receiver, 421 _ => panic!("follower awaits"), 422 }; 423 leader.too_large(); 424 assert!( 425 matches!(wait(receiver).await, Resolved::Regenerate), 426 "an oversized leader sends its followers to stream directly" 427 ); 428 match cache.decide(key(b"huge")) { 429 Decision::Stream => {} 430 _ => panic!("an oversized key streams directly without re-buffering"), 431 } 432 } 433 434 #[tokio::test] 435 async fn an_expired_entry_is_regenerated() { 436 let clock = Arc::new(ManualClock::new(UnixMicros::new(0))); 437 let cache = PackCache::new( 438 CacheConfig { 439 ttl: Duration::from_secs(60), 440 ..CacheConfig::default() 441 }, 442 Arc::clone(&clock) as Arc<dyn Clock>, 443 ); 444 lead(&cache, b"stale").ready(Bytes::from_static(b"old")); 445 clock.advance(Duration::from_secs(61)); 446 match cache.decide(key(b"stale")) { 447 Decision::Lead(_) => {} 448 _ => panic!("an expired entry forces a fresh generation"), 449 } 450 } 451 452 #[tokio::test] 453 async fn the_total_byte_limit_evicts_the_oldest_entry() { 454 let cache = built(CacheConfig { 455 max_total_bytes: MaxCacheBytes::new(8), 456 ..CacheConfig::default() 457 }); 458 lead(&cache, b"a").ready(Bytes::from(vec![0u8; 5])); 459 lead(&cache, b"b").ready(Bytes::from(vec![0u8; 5])); 460 match cache.decide(key(b"a")) { 461 Decision::Lead(_) => {} 462 _ => panic!("the oldest entry is evicted once the total limit is exceeded"), 463 } 464 match cache.decide(key(b"b")) { 465 Decision::Serve(_) => {} 466 _ => panic!("the newest entry survives eviction"), 467 } 468 } 469 470 #[tokio::test] 471 async fn the_entry_limit_never_exceeds_the_total_cache_size() { 472 let cache = built(CacheConfig { 473 max_entry_bytes: MaxEntryBytes::new(64), 474 max_total_bytes: MaxCacheBytes::new(16), 475 ..CacheConfig::default() 476 }); 477 let lease = lead(&cache, b"probe"); 478 assert_eq!( 479 lease.max_entry_bytes().get(), 480 16, 481 "a per-entry limit above the whole-cache size is clamped so a full entry can be retained" 482 ); 483 } 484 485 #[tokio::test] 486 async fn oversized_entries_cannot_grow_the_cache_without_bound() { 487 let cache = built(CacheConfig::default()); 488 (0..(MAX_ENTRIES + 64)).for_each(|nonce| { 489 lead(&cache, format!("oversized-{nonce}").as_bytes()).too_large(); 490 }); 491 assert!( 492 cache.entry_count() <= MAX_ENTRIES as u64, 493 "a flood of distinct oversized requests stays within the entry bound" 494 ); 495 } 496 497 #[test] 498 fn capture_stops_buffering_once_the_limit_is_passed() { 499 let mut capture = Capture::new(Some(MaxEntryBytes::new(4))); 500 capture.record(b"abc"); 501 capture.record(b"de"); 502 assert!(capture.into_bytes().is_none(), "overflow drops the buffer"); 503 } 504 505 #[test] 506 fn capture_keeps_bytes_under_the_limit() { 507 let mut capture = Capture::new(Some(MaxEntryBytes::new(8))); 508 capture.record(b"abcd"); 509 assert_eq!(capture.into_bytes().unwrap(), b"abcd"); 510 } 511}