Coverage Report

Created: 2026-07-21 15:28

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/build/source/nativelink-store/src/ref_store.rs
Line
Count
Source
1
// Copyright 2024 The NativeLink Authors. All rights reserved.
2
//
3
// Licensed under the Functional Source License, Version 1.1, Apache 2.0 Future License (the "License");
4
// you may not use this file except in compliance with the License.
5
// You may obtain a copy of the License at
6
//
7
//    See LICENSE file for details
8
//
9
// Unless required by applicable law or agreed to in writing, software
10
// distributed under the License is distributed on an "AS IS" BASIS,
11
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12
// See the License for the specific language governing permissions and
13
// limitations under the License.
14
15
use core::cell::UnsafeCell;
16
use core::pin::Pin;
17
use std::sync::{Arc, Weak};
18
19
use async_trait::async_trait;
20
use nativelink_config::stores::RefSpec;
21
use nativelink_error::{Error, ResultExt, make_input_err};
22
use nativelink_metric::MetricsComponent;
23
use nativelink_util::buf_channel::{DropCloserReadHalf, DropCloserWriteHalf};
24
use nativelink_util::health_utils::{HealthStatusIndicator, default_health_status_indicator};
25
use nativelink_util::store_trait::{
26
    RemoveCallback, Store, StoreDriver, StoreKey, StoreLike, UploadSizeInfo,
27
};
28
use parking_lot::Mutex;
29
use tracing::{debug, error};
30
31
use crate::store_manager::StoreManager;
32
33
#[repr(C, align(8))]
34
#[derive(Debug)]
35
struct AlignedStoreCell(UnsafeCell<Option<Store>>);
36
37
#[derive(Debug)]
38
struct StoreReference {
39
    cell: AlignedStoreCell,
40
    mux: Mutex<()>,
41
}
42
43
unsafe impl Sync for StoreReference {}
44
45
#[derive(Debug, MetricsComponent)]
46
pub struct RefStore {
47
    #[metric(help = "The store we are referencing")]
48
    name: String,
49
    store_manager: Weak<StoreManager>,
50
    inner: StoreReference,
51
    remove_callbacks: Mutex<Vec<RemoveCallback>>,
52
}
53
54
impl RefStore {
55
6
    pub fn new(spec: &RefSpec, store_manager: Weak<StoreManager>) -> Arc<Self> {
56
6
        Arc::new(Self {
57
6
            name: spec.name.clone(),
58
6
            store_manager,
59
6
            inner: StoreReference {
60
6
                mux: Mutex::new(()),
61
6
                cell: AlignedStoreCell(UnsafeCell::new(None)),
62
6
            },
63
6
            remove_callbacks: Mutex::new(vec![]),
64
6
        })
65
6
    }
66
67
    #[inline]
68
6
    fn get_store(&self) -> Result<&Store, Error> {
69
6
        let ref_store = self.inner.cell.0.get();
70
        unsafe {
71
6
            if let Some(
ref store5
) = *ref_store {
72
5
                return Ok(store);
73
1
            }
74
        }
75
1
        Err(make_input_err!(
76
1
            "ref_store cannot get store '{}', was post_init called?",
77
1
            self.name
78
1
        ))
79
6
    }
80
}
81
82
#[async_trait]
83
impl StoreDriver for RefStore {
84
5
    async fn post_init(self: Arc<Self>) -> Result<(), Error> {
85
        debug!("Running post_init to get store");
86
        let ref_store = self.inner.cell.0.get();
87
        unsafe {
88
            if (*ref_store).is_some() {
89
                return Err(make_input_err!("post_init already called for RefStore"));
90
            }
91
        }
92
        let _lock = self.inner.mux.lock();
93
        let store_manager = self
94
            .store_manager
95
            .upgrade()
96
            .err_tip(|| "Store manager is gone")?;
97
        if let Some(store) = store_manager.get_store(&self.name) {
98
            let remove_callbacks = self.remove_callbacks.lock().clone();
99
            for callback in remove_callbacks {
100
                debug!(?callback, ?store, "Added callback to store");
101
                store.register_remove_callback(callback)?;
102
            }
103
            unsafe {
104
                *ref_store = Some(store);
105
            }
106
            Ok(())
107
        } else {
108
            Err(make_input_err!(
109
                "Failed to find store '{}' in StoreManager in RefStore",
110
                self.name
111
            ))
112
        }
113
5
    }
114
115
    async fn has_with_results(
116
        self: Pin<&Self>,
117
        keys: &[StoreKey<'_>],
118
        results: &mut [Option<u64>],
119
2
    ) -> Result<(), Error> {
120
        self.get_store()?.has_with_results(keys, results).await
121
2
    }
122
123
    async fn update(
124
        self: Pin<&Self>,
125
        key: StoreKey<'_>,
126
        reader: DropCloserReadHalf,
127
        size_info: UploadSizeInfo,
128
1
    ) -> Result<u64, Error> {
129
        self.get_store()?.update(key, reader, size_info).await
130
1
    }
131
132
    async fn get_part(
133
        self: Pin<&Self>,
134
        key: StoreKey<'_>,
135
        writer: &mut DropCloserWriteHalf,
136
        offset: u64,
137
        length: Option<u64>,
138
1
    ) -> Result<(), Error> {
139
        self.get_store()?
140
            .get_part(key, writer, offset, length)
141
            .await
142
1
    }
143
144
2
    fn inner_store(&self, key: Option<StoreKey>) -> &'_ dyn StoreDriver {
145
2
        match self.get_store() {
146
2
            Ok(store) => store.inner_store(key),
147
0
            Err(err) => {
148
0
                error!(?key, ?err, "Failed to get store for key");
149
0
                self
150
            }
151
        }
152
2
    }
153
154
0
    fn as_any<'a>(&'a self) -> &'a (dyn core::any::Any + Sync + Send + 'static) {
155
0
        self
156
0
    }
157
158
0
    fn as_any_arc(self: Arc<Self>) -> Arc<dyn core::any::Any + Sync + Send + 'static> {
159
0
        self
160
0
    }
161
162
0
    fn register_remove_callback(self: Arc<Self>, callback: RemoveCallback) -> Result<(), Error> {
163
0
        self.remove_callbacks.lock().push(callback.clone());
164
0
        let ref_store = self.inner.cell.0.get();
165
        unsafe {
166
0
            if let Some(ref store) = *ref_store {
167
0
                debug!(?callback, ?store, "New callback");
168
0
                store.register_remove_callback(callback)?;
169
            } else {
170
0
                debug!(?callback, "New callback, no store yet");
171
            }
172
        }
173
0
        Ok(())
174
0
    }
175
}
176
177
default_health_status_indicator!(RefStore);