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 "fuse.gcsfuse",
27 "fuse.juicefs",
28 "davfs",
29 "ftpfs",
30 "autofs",
34];
35
36const MEMORY_FILESYSTEMS: &[&str] = &["tmpfs", "ramfs", "devtmpfs"];
39
40#[derive(Debug, Clone, Copy, PartialEq, Eq)]
42pub enum Locality {
43 Local,
45 Memory,
47 Network,
49 Object,
51 Unknown,
53}
54
55impl Locality {
56 pub fn of_fstype(fstype: &str) -> Locality {
60 match fstype {
61 "s3" | "s3a" | "gs" | "gcs" | "az" | "cloud" | "http" | "https" => Locality::Object,
62 "" | "unknown" => Locality::Unknown,
63 other => classify(other),
64 }
65 }
66}
67
68#[derive(Debug, Clone, PartialEq, Eq)]
70pub struct Source {
71 pub fstype: String,
74 pub locality: Locality,
75}
76
77impl Source {
78 pub fn label(&self) -> &str {
82 &self.fstype
83 }
84
85 pub fn network(&self) -> bool {
86 self.locality == Locality::Network
87 }
88
89 pub fn from_fstype(fstype: &str) -> Self {
92 if let Some(scheme) = ["s3", "gs", "http", "https", "az", "hdfs"]
93 .into_iter()
94 .find(|s| *s == fstype)
95 {
96 return Self {
97 fstype: scheme.to_string(),
98 locality: Locality::Object,
99 };
100 }
101 Self {
102 locality: classify(fstype),
103 fstype: fstype.to_string(),
104 }
105 }
106}
107
108const MOUNTS_TTL: Duration = Duration::from_secs(5);
112
113static CACHED_MOUNTS: Mutex<Option<(Instant, Arc<Mounts>)>> = Mutex::new(None);
115
116#[derive(Debug, Clone, Default)]
119pub struct Mounts {
120 entries: Vec<(String, String)>,
122}
123
124impl Mounts {
125 pub fn current() -> Self {
128 std::fs::read_to_string("/proc/self/mountinfo")
129 .map(|s| Self::parse(&s))
130 .unwrap_or_default()
131 }
132
133 pub fn cached() -> Arc<Self> {
136 let now = Instant::now();
137 let mut slot = CACHED_MOUNTS
141 .lock()
142 .unwrap_or_else(std::sync::PoisonError::into_inner);
143 if let Some((read_at, mounts)) = slot.as_ref()
144 && now.duration_since(*read_at) < MOUNTS_TTL
145 {
146 return Arc::clone(mounts);
147 }
148 let mounts = Arc::new(Self::current());
149 *slot = Some((now, Arc::clone(&mounts)));
150 mounts
151 }
152
153 pub fn parse(mountinfo: &str) -> Self {
154 let mut entries = Vec::new();
155 for line in mountinfo.lines() {
156 let Some((before, after)) = line.split_once(" - ") else {
159 continue;
160 };
161 let Some(point) = before.split_whitespace().nth(4) else {
162 continue;
163 };
164 let Some(fstype) = after.split_whitespace().next() else {
165 continue;
166 };
167 entries.push((point.to_string(), fstype.to_string()));
168 }
169 Self { entries }
170 }
171
172 pub fn fstype_for(&self, path: &Path) -> Option<&str> {
175 self.covering(path).map(|(_, fstype)| fstype)
176 }
177
178 pub fn mount_point_for(&self, path: &Path) -> std::path::PathBuf {
180 self.covering(path)
181 .map(|(point, _)| std::path::PathBuf::from(point))
182 .unwrap_or_default()
183 }
184
185 fn covering(&self, path: &Path) -> Option<(&str, &str)> {
187 let joined;
190 let path = if path.has_root() {
193 path
194 } else {
195 match std::env::current_dir() {
196 Ok(cwd) => {
197 joined = cwd.join(path);
198 &joined
199 }
200 Err(_) => path,
201 }
202 };
203
204 let mut best: Option<(&str, &str)> = None;
205 for (point, fstype) in &self.entries {
206 if !path.starts_with(point) {
207 continue;
208 }
209 if best.is_none_or(|(at, _)| point.len() >= at.len()) {
210 best = Some((point, fstype));
211 }
212 }
213 best
214 }
215
216 pub fn describe(&self, path: &Path) -> Source {
218 if let Some(scheme) = object_scheme(path) {
219 return Source {
220 fstype: scheme,
221 locality: Locality::Object,
222 };
223 }
224 match self.fstype_for(path) {
225 Some(fstype) => Source {
226 locality: classify(fstype),
227 fstype: fstype.to_string(),
228 },
229 None => Source {
230 fstype: "unknown".to_string(),
231 locality: Locality::Unknown,
232 },
233 }
234 }
235
236 pub fn is_network(&self, path: &Path) -> bool {
237 self.describe(path).network()
238 }
239
240 pub fn could_block(&self, path: &Path) -> bool {
245 let source = self.describe(path);
246 source.network() || source.fstype.starts_with("fuse.")
247 }
248}
249
250fn classify(fstype: &str) -> Locality {
251 if NETWORK_FILESYSTEMS.contains(&fstype) {
252 Locality::Network
253 } else if MEMORY_FILESYSTEMS.contains(&fstype) {
254 Locality::Memory
255 } else {
256 Locality::Local
257 }
258}
259
260pub fn object_scheme(path: &Path) -> Option<String> {
263 match crate::cloud::source::input_source(path) {
264 crate::cloud::source::InputSource::Local(_) => None,
265 crate::cloud::source::InputSource::S3(_) => Some("s3".to_string()),
266 crate::cloud::source::InputSource::Gcs(_) => Some("gs".to_string()),
267 crate::cloud::source::InputSource::Azure(_) => Some("az".to_string()),
268 crate::cloud::source::InputSource::Http(_) => Some("http".to_string()),
269 }
270}
271
272#[cfg(test)]
273mod locality_of_fstype_tests {
274 use super::*;
275
276 #[test]
277 fn object_store_schemes_are_not_local_disks() {
278 for scheme in ["s3", "s3a", "gs", "gcs", "http", "https"] {
282 assert_eq!(
283 Locality::of_fstype(scheme),
284 Locality::Object,
285 "{scheme} should be an object store"
286 );
287 }
288 }
289
290 #[test]
291 fn network_filesystems_are_network() {
292 for fstype in ["nfs", "nfs4", "cifs", "smb3"] {
293 assert_eq!(
294 Locality::of_fstype(fstype),
295 Locality::Network,
296 "{fstype} should be network"
297 );
298 }
299 }
300
301 #[test]
302 fn memory_filesystems_are_memory() {
303 assert_eq!(Locality::of_fstype("tmpfs"), Locality::Memory);
304 }
305
306 #[test]
307 fn ordinary_filesystems_are_local() {
308 for fstype in ["ext4", "btrfs", "xfs", "apfs", "ntfs"] {
309 assert_eq!(
310 Locality::of_fstype(fstype),
311 Locality::Local,
312 "{fstype} should be local"
313 );
314 }
315 }
316
317 #[test]
318 fn nothing_known_is_not_guessed_at() {
319 assert_eq!(Locality::of_fstype(""), Locality::Unknown);
320 assert_eq!(Locality::of_fstype("unknown"), Locality::Unknown);
321 }
322
323 #[test]
326 fn it_agrees_with_describe() {
327 let mounts = Mounts::parse("");
328 for path in ["s3://bucket/key.parquet", "gs://bucket/key.parquet"] {
329 let source = mounts.describe(std::path::Path::new(path));
330 assert_eq!(source.locality, Locality::of_fstype(&source.fstype));
331 }
332 }
333}