This repository has no description
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}