1use super::*;
4use crate::async_lru_cache::AsyncLruCacheBackend;
5use tracing::trace;
6
7struct L2CacheBackend<S: Storage> {
9 file: Arc<S>,
11
12 header: Arc<Header>,
14}
15
16struct RefBlockCacheBackend<S: Storage> {
18 file: Arc<S>,
20
21 header: Arc<Header>,
23}
24
25impl<S: Storage> L2CacheBackend<S> {
26 pub fn new(file: Arc<S>, header: Arc<Header>) -> Self {
30 L2CacheBackend { file, header }
31 }
32}
33
34#[maybe_async(AFIT)]
35impl<S: Storage> AsyncLruCacheBackend for L2CacheBackend<S> {
36 type Key = HostCluster;
37 type Value = L2Table;
38
39 async fn load(&self, l2_cluster: HostCluster) -> io::Result<L2Table> {
40 trace!("Loading L2 table");
41
42 L2Table::load(
43 self.file.as_ref(),
44 &self.header,
45 l2_cluster,
46 self.header.l2_entries(),
47 )
48 .await
49 }
50
51 async fn flush(&self, l2_cluster: HostCluster, l2_table: &L2Table) -> io::Result<()> {
52 trace!("Flushing L2 table");
53 if l2_table.is_modified() {
54 assert!(l2_table.get_cluster().unwrap() == l2_cluster);
55 l2_table.write(self.file.as_ref()).await?;
56 }
57 Ok(())
58 }
59
60 unsafe fn evict(&self, _l2_cluster: HostCluster, l2_table: L2Table) {
61 trace!(
62 "Evicting L2 table {}",
63 l2_table.get_offset().unwrap_or(HostOffset(0))
64 );
65 l2_table.clear_modified();
66 }
67}
68
69impl<S: Storage> RefBlockCacheBackend<S> {
70 pub fn new(file: Arc<S>, header: Arc<Header>) -> Self {
74 RefBlockCacheBackend { file, header }
75 }
76}
77
78#[maybe_async(AFIT)]
79impl<S: Storage> AsyncLruCacheBackend for RefBlockCacheBackend<S> {
80 type Key = HostCluster;
81 type Value = RefBlock;
82
83 async fn load(&self, rb_cluster: HostCluster) -> io::Result<RefBlock> {
84 RefBlock::load(self.file.as_ref(), &self.header, rb_cluster).await
85 }
86
87 async fn flush(&self, rb_cluster: HostCluster, refblock: &RefBlock) -> io::Result<()> {
88 if refblock.is_modified() {
89 assert!(refblock.get_cluster().unwrap() == rb_cluster);
90 refblock.write(self.file.as_ref()).await?;
91 }
92 Ok(())
93 }
94
95 unsafe fn evict(&self, _rb_cluster: HostCluster, refblock: RefBlock) {
96 refblock.clear_modified();
97 }
98}
99
100#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
102enum CacheDependency {
103 #[default]
105 None,
106 L2DependsOnRb,
108 RbDependsOnL2,
110}
111
112pub(super) struct MetadataCaches<S: Storage> {
126 l2: AsyncLruCache<HostCluster, L2Table, L2CacheBackend<S>>,
128
129 rb: AsyncLruCache<HostCluster, RefBlock, RefBlockCacheBackend<S>>,
131
132 direction: RwLock<CacheDependency>,
138}
139
140#[maybe_async]
141impl<S: Storage> MetadataCaches<S> {
142 pub fn new(file: &Arc<S>, header: &Arc<Header>, l2_entries: usize, rb_entries: usize) -> Self {
147 let l2_backend = L2CacheBackend::new(Arc::clone(file), Arc::clone(header));
148 let rb_backend = RefBlockCacheBackend::new(Arc::clone(file), Arc::clone(header));
149
150 MetadataCaches {
151 l2: AsyncLruCache::new(l2_backend, l2_entries),
152 rb: AsyncLruCache::new(rb_backend, rb_entries),
153 direction: Default::default(),
154 }
155 }
156
157 pub async fn l2_depends_on_rb(&self) -> io::Result<()> {
161 let mut dir = self.direction.write().await;
162 if *dir == CacheDependency::L2DependsOnRb {
163 return Ok(());
164 }
165 if *dir == CacheDependency::RbDependsOnL2 {
166 self.l2.flush().await?;
167 }
168 *dir = CacheDependency::L2DependsOnRb;
169 Ok(())
170 }
171
172 pub async fn rb_depends_on_l2(&self) -> io::Result<()> {
176 let mut dir = self.direction.write().await;
177 if *dir == CacheDependency::RbDependsOnL2 {
178 return Ok(());
179 }
180 if *dir == CacheDependency::L2DependsOnRb {
181 self.rb.flush().await?;
182 }
183 *dir = CacheDependency::RbDependsOnL2;
184 Ok(())
185 }
186
187 pub async fn flush_all(&self) -> io::Result<()> {
189 let dir = self.direction.read().await;
190 if *dir == CacheDependency::L2DependsOnRb {
191 self.rb.flush().await?;
192 self.l2.flush().await?;
193 } else {
194 self.l2.flush().await?;
195 self.rb.flush().await?;
196 }
197
198 Ok(())
199 }
200
201 pub async fn l2_get_or_insert(&self, cluster_index: HostCluster) -> io::Result<Arc<L2Table>> {
205 let dir = self.direction.read().await;
206 if let Some(l2) = self.l2.get_or_insert(cluster_index, false).await? {
207 return Ok(l2);
208 }
209
210 if *dir == CacheDependency::L2DependsOnRb {
211 self.rb.flush().await?;
212 }
213
214 let l2 = self.l2.get_or_insert(cluster_index, true).await?.unwrap();
216 Ok(l2)
217 }
218
219 pub async fn l2_insert(
223 &self,
224 cluster_index: HostCluster,
225 table: Arc<L2Table>,
226 ) -> io::Result<()> {
227 let dir = self.direction.read().await;
228 if !self
229 .l2
230 .insert(cluster_index, Arc::clone(&table), false)
231 .await?
232 {
233 if *dir == CacheDependency::L2DependsOnRb {
234 self.rb.flush().await?;
235 }
236
237 let inserted = self.l2.insert(cluster_index, table, true).await?;
238 assert!(inserted);
240 }
241
242 Ok(())
243 }
244
245 pub async unsafe fn invalidate_l2(&self) -> io::Result<()> {
250 unsafe { self.l2.invalidate() }.await
251 }
252
253 pub async fn rb_get_or_insert(&self, cluster_index: HostCluster) -> io::Result<Arc<RefBlock>> {
257 let dir = self.direction.read().await;
258 if let Some(rb) = self.rb.get_or_insert(cluster_index, false).await? {
259 return Ok(rb);
260 }
261
262 if *dir == CacheDependency::RbDependsOnL2 {
263 self.l2.flush().await?;
264 }
265
266 let rb = self.rb.get_or_insert(cluster_index, true).await?.unwrap();
268 Ok(rb)
269 }
270
271 pub async fn rb_insert(&self, cluster_index: HostCluster, rb: Arc<RefBlock>) -> io::Result<()> {
275 let dir = self.direction.read().await;
276 if !self
277 .rb
278 .insert(cluster_index, Arc::clone(&rb), false)
279 .await?
280 {
281 if *dir == CacheDependency::RbDependsOnL2 {
282 self.l2.flush().await?;
283 }
284
285 let inserted = self.rb.insert(cluster_index, rb, true).await?;
286 assert!(inserted);
288 }
289
290 Ok(())
291 }
292
293 pub async fn flush_rb(&self) -> io::Result<()> {
295 let dir = self.direction.read().await;
296 if *dir == CacheDependency::RbDependsOnL2 {
297 self.l2.flush().await?;
298 }
299 self.rb.flush().await
300 }
301
302 pub async unsafe fn invalidate_rb(&self) -> io::Result<()> {
307 unsafe { self.rb.invalidate() }.await
308 }
309}
310
311#[cfg(test)]
312mod tests {
313 use super::*;
314 use crate::null::Null;
315
316 fn make_test_caches() -> MetadataCaches<Null> {
317 let null = Arc::new(Null::new(1 << 16));
318 let header = Arc::new(Header::new(16, 1, None, None, None));
319 MetadataCaches::new(&null, &header, 16, 16)
320 }
321
322 #[cfg(feature = "sync")]
323 fn block_on<T>(v: T) -> T {
324 v
325 }
326
327 #[cfg(feature = "async")]
328 fn block_on<F: std::future::Future>(f: F) -> F::Output {
329 tokio::runtime::Builder::new_current_thread()
330 .build()
331 .unwrap()
332 .block_on(f)
333 }
334
335 #[maybe_async::test(feature = "sync", async(feature = "async", tokio::test))]
337 async fn test_direction_switch() {
338 let caches = make_test_caches();
339
340 caches.l2_depends_on_rb().await.unwrap();
341 caches.l2_depends_on_rb().await.unwrap();
342 caches.rb_depends_on_l2().await.unwrap();
343 caches.rb_depends_on_l2().await.unwrap();
344 caches.l2_depends_on_rb().await.unwrap();
345 }
346
347 #[test]
354 fn test_cross_cache_eviction_no_deadlock() {
355 use std::thread;
356
357 let null = Arc::new(Null::new(1 << 20));
358 let header = Arc::new(Header::new(16, 1, None, None, None));
359 let caches = Arc::new(MetadataCaches::new(&null, &header, 1, 1));
360
361 block_on(caches.l2_depends_on_rb()).unwrap();
362
363 let mut l2_entry = L2Table::new_cleared(&header);
366 l2_entry.set_cluster(HostCluster(0));
367 l2_entry.clear_modified();
368 let mut rb_entry = RefBlock::new_cleared(null.as_ref(), &header).unwrap();
369 rb_entry.set_cluster(HostCluster(0));
370 rb_entry.clear_modified();
371
372 block_on(caches.l2_insert(HostCluster(0), Arc::new(l2_entry))).unwrap();
373 block_on(caches.rb_insert(HostCluster(0), Arc::new(rb_entry))).unwrap();
374
375 let c1 = Arc::clone(&caches);
376 let c2 = Arc::clone(&caches);
377
378 let barrier = Arc::new(std::sync::Barrier::new(2));
379 let bar1 = Arc::clone(&barrier);
380 let bar2 = Arc::clone(&barrier);
381
382 let header1 = Arc::clone(&header);
384 let t1 = thread::spawn(move || {
385 let mut entry = L2Table::new_cleared(&header1);
386 entry.set_cluster(HostCluster(1));
387 entry.clear_modified();
388 bar1.wait();
389 block_on(c1.l2_insert(HostCluster(1), Arc::new(entry))).unwrap();
390 });
391
392 let header2 = Arc::clone(&header);
394 let null2 = Arc::new(Null::new(1 << 20));
395 let t2 = thread::spawn(move || {
396 let mut entry = RefBlock::new_cleared(null2.as_ref(), &header2).unwrap();
397 entry.set_cluster(HostCluster(1));
398 entry.clear_modified();
399 bar2.wait();
400 block_on(c2.rb_insert(HostCluster(1), Arc::new(entry))).unwrap();
401 });
402
403 t1.join().unwrap();
404 t2.join().unwrap();
405 }
406
407 #[test]
410 fn test_concurrent_direction_switch() {
411 use std::thread;
412
413 let caches = Arc::new(make_test_caches());
414
415 let c1 = Arc::clone(&caches);
416 let c2 = Arc::clone(&caches);
417
418 let barrier = Arc::new(std::sync::Barrier::new(2));
419 let bar1 = Arc::clone(&barrier);
420 let bar2 = Arc::clone(&barrier);
421
422 let t1 = thread::spawn(move || {
423 bar1.wait();
424 block_on(c1.l2_depends_on_rb()).unwrap();
425 });
426
427 let t2 = thread::spawn(move || {
428 bar2.wait();
429 block_on(c2.rb_depends_on_l2()).unwrap();
430 });
431
432 t1.join().unwrap();
433 t2.join().unwrap();
434 }
435
436 #[maybe_async::test(feature = "sync", async(feature = "async", tokio::test))]
438 async fn test_direction_switch_flushes_opposing() {
439 let null = Arc::new(Null::new(1 << 20));
440 let header = Arc::new(Header::new(16, 1, None, None, None));
441 let caches = MetadataCaches::new(&null, &header, 16, 16);
442
443 let mut rb_entry = RefBlock::new_cleared(null.as_ref(), &header).unwrap();
445 rb_entry.set_cluster(HostCluster(0));
446 let rb_arc = Arc::new(rb_entry);
448 caches
449 .rb_insert(HostCluster(0), Arc::clone(&rb_arc))
450 .await
451 .unwrap();
452
453 assert!(rb_arc.is_modified());
454
455 caches.l2_depends_on_rb().await.unwrap();
457 caches.rb_depends_on_l2().await.unwrap();
458
459 assert!(!rb_arc.is_modified());
461 }
462
463 #[maybe_async::test(feature = "sync", async(feature = "async", tokio::test))]
468 async fn test_dep_flushed_before_eviction() {
469 let null = Arc::new(Null::new(1 << 20));
470 let header = Arc::new(Header::new(16, 1, None, None, None));
471 let caches = MetadataCaches::new(&null, &header, 1, 16);
472
473 caches.l2_depends_on_rb().await.unwrap();
474
475 let mut rb_entry = RefBlock::new_cleared(null.as_ref(), &header).unwrap();
477 rb_entry.set_cluster(HostCluster(0));
478 let rb_arc = Arc::new(rb_entry);
479 caches
480 .rb_insert(HostCluster(0), Arc::clone(&rb_arc))
481 .await
482 .unwrap();
483
484 let mut l2_entry = L2Table::new_cleared(&header);
486 l2_entry.set_cluster(HostCluster(0));
487 l2_entry.clear_modified();
488 caches
489 .l2_insert(HostCluster(0), Arc::new(l2_entry))
490 .await
491 .unwrap();
492
493 let mut l2_entry2 = L2Table::new_cleared(&header);
496 l2_entry2.set_cluster(HostCluster(1));
497 l2_entry2.clear_modified();
498 caches
499 .l2_insert(HostCluster(1), Arc::new(l2_entry2))
500 .await
501 .unwrap();
502
503 assert!(
505 !rb_arc.is_modified(),
506 "rb entry must be flushed before l2 eviction"
507 );
508 }
509}