1use std::path::Path;
7use std::sync::{Arc, Mutex};
8use std::time::{Duration, Instant};
9
10pub const NETWORK_FILESYSTEMS: &[&str] = &[
13 "nfs",
14 "nfs4",
15 "cifs",
16 "smb3",
17 "smbfs",
18 "afs",
19 "9p",
20 "ceph",
21 "glusterfs",
22 "fuse.sshfs",
23 "fuse.rclone",
24 "fuse.s3fs",
25 "fuse.davfs",
26 "davfs",
27 "ftpfs",
28 "autofs",
32];
33
34const MEMORY_FILESYSTEMS: &[&str] = &["tmpfs", "ramfs", "devtmpfs"];
37
38#[derive(Debug, Clone, Copy, PartialEq, Eq)]
40pub enum Locality {
41 Local,
43 Memory,
45 Network,
47 Object,
49 Unknown,
51}
52
53impl Locality {
54 pub fn of_fstype(fstype: &str) -> Locality {
58 match fstype {
59 "s3" | "s3a" | "gs" | "gcs" | "az" | "cloud" | "http" | "https" => Locality::Object,
60 "" | "unknown" => Locality::Unknown,
61 other => classify(other),
62 }
63 }
64}
65
66#[derive(Debug, Clone, PartialEq, Eq)]
68pub struct Source {
69 pub fstype: String,
72 pub locality: Locality,
73}
74
75impl Source {
76 pub fn label(&self) -> &str {
80 &self.fstype
81 }
82
83 pub fn network(&self) -> bool {
84 self.locality == Locality::Network
85 }
86
87 pub fn from_fstype(fstype: &str) -> Self {
90 if let Some(scheme) = ["s3", "gs", "http", "https", "az", "hdfs"]
91 .into_iter()
92 .find(|s| *s == fstype)
93 {
94 return Self {
95 fstype: scheme.to_string(),
96 locality: Locality::Object,
97 };
98 }
99 Self {
100 locality: classify(fstype),
101 fstype: fstype.to_string(),
102 }
103 }
104}
105
106const MOUNTS_TTL: Duration = Duration::from_millis(500);
110
111static CACHED_MOUNTS: Mutex<Option<(Instant, Arc<Mounts>)>> = Mutex::new(None);
113
114#[derive(Debug, Clone, Default)]
117pub struct Mounts {
118 entries: Vec<(String, String)>,
120}
121
122impl Mounts {
123 pub fn current() -> Self {
126 std::fs::read_to_string("/proc/self/mountinfo")
127 .map(|s| Self::parse(&s))
128 .unwrap_or_default()
129 }
130
131 pub fn cached() -> Arc<Self> {
134 let now = Instant::now();
135 let mut slot = CACHED_MOUNTS
139 .lock()
140 .unwrap_or_else(std::sync::PoisonError::into_inner);
141 if let Some((read_at, mounts)) = slot.as_ref()
142 && now.duration_since(*read_at) < MOUNTS_TTL
143 {
144 return Arc::clone(mounts);
145 }
146 let mounts = Arc::new(Self::current());
147 *slot = Some((now, Arc::clone(&mounts)));
148 mounts
149 }
150
151 pub fn parse(mountinfo: &str) -> Self {
152 let mut entries = Vec::new();
153 for line in mountinfo.lines() {
154 let Some((before, after)) = line.split_once(" - ") else {
157 continue;
158 };
159 let Some(point) = before.split_whitespace().nth(4) else {
160 continue;
161 };
162 let Some(fstype) = after.split_whitespace().next() else {
163 continue;
164 };
165 entries.push((point.to_string(), fstype.to_string()));
166 }
167 Self { entries }
168 }
169
170 pub fn fstype_for(&self, path: &Path) -> Option<&str> {
173 let joined;
176 let path = if path.has_root() {
179 path
180 } else {
181 match std::env::current_dir() {
182 Ok(cwd) => {
183 joined = cwd.join(path);
184 &joined
185 }
186 Err(_) => path,
187 }
188 };
189
190 let mut best: Option<(usize, &str)> = None;
191 for (point, fstype) in &self.entries {
192 if !path.starts_with(point) {
193 continue;
194 }
195 let len = point.len();
196 if best.is_none_or(|(n, _)| len >= n) {
197 best = Some((len, fstype));
198 }
199 }
200 best.map(|(_, f)| f)
201 }
202
203 pub fn describe(&self, path: &Path) -> Source {
205 if let Some(scheme) = object_scheme(path) {
206 return Source {
207 fstype: scheme,
208 locality: Locality::Object,
209 };
210 }
211 match self.fstype_for(path) {
212 Some(fstype) => Source {
213 locality: classify(fstype),
214 fstype: fstype.to_string(),
215 },
216 None => Source {
217 fstype: "unknown".to_string(),
218 locality: Locality::Unknown,
219 },
220 }
221 }
222
223 pub fn is_network(&self, path: &Path) -> bool {
224 self.describe(path).network()
225 }
226}
227
228fn classify(fstype: &str) -> Locality {
229 if NETWORK_FILESYSTEMS.contains(&fstype) {
230 Locality::Network
231 } else if MEMORY_FILESYSTEMS.contains(&fstype) {
232 Locality::Memory
233 } else {
234 Locality::Local
235 }
236}
237
238pub fn object_scheme(path: &Path) -> Option<String> {
241 match crate::cloud::source::input_source(path) {
242 crate::cloud::source::InputSource::Local(_) => None,
243 crate::cloud::source::InputSource::S3(_) => Some("s3".to_string()),
244 crate::cloud::source::InputSource::Gcs(_) => Some("gs".to_string()),
245 crate::cloud::source::InputSource::Azure(_) => Some("az".to_string()),
246 crate::cloud::source::InputSource::Http(_) => Some("http".to_string()),
247 }
248}
249
250#[cfg(test)]
251mod locality_of_fstype_tests {
252 use super::*;
253
254 #[test]
255 fn object_store_schemes_are_not_local_disks() {
256 for scheme in ["s3", "s3a", "gs", "gcs", "http", "https"] {
260 assert_eq!(
261 Locality::of_fstype(scheme),
262 Locality::Object,
263 "{scheme} should be an object store"
264 );
265 }
266 }
267
268 #[test]
269 fn network_filesystems_are_network() {
270 for fstype in ["nfs", "nfs4", "cifs", "smb3"] {
271 assert_eq!(
272 Locality::of_fstype(fstype),
273 Locality::Network,
274 "{fstype} should be network"
275 );
276 }
277 }
278
279 #[test]
280 fn memory_filesystems_are_memory() {
281 assert_eq!(Locality::of_fstype("tmpfs"), Locality::Memory);
282 }
283
284 #[test]
285 fn ordinary_filesystems_are_local() {
286 for fstype in ["ext4", "btrfs", "xfs", "apfs", "ntfs"] {
287 assert_eq!(
288 Locality::of_fstype(fstype),
289 Locality::Local,
290 "{fstype} should be local"
291 );
292 }
293 }
294
295 #[test]
296 fn nothing_known_is_not_guessed_at() {
297 assert_eq!(Locality::of_fstype(""), Locality::Unknown);
298 assert_eq!(Locality::of_fstype("unknown"), Locality::Unknown);
299 }
300
301 #[test]
304 fn it_agrees_with_describe() {
305 let mounts = Mounts::parse("");
306 for path in ["s3://bucket/key.parquet", "gs://bucket/key.parquet"] {
307 let source = mounts.describe(std::path::Path::new(path));
308 assert_eq!(source.locality, Locality::of_fstype(&source.fstype));
309 }
310 }
311}