1use std::borrow::Cow;
2
3use futures::TryStreamExt;
4use object_store::path::Path;
5use polars_error::{PolarsResult, polars_bail, polars_err};
6use polars_utils::pl_path::{CloudScheme, PlRefPath};
7use polars_utils::pl_str::PlSmallStr;
8use regex::Regex;
9
10use super::CloudOptions;
11
12pub(crate) fn extract_prefix_expansion(path: &str) -> PolarsResult<(Cow<'_, str>, Option<String>)> {
18 let mut replacements: Vec<(usize, usize, &[u8])> = vec![];
20
21 let mut pos: usize = if let Some(after_last_slash) =
25 memchr::memchr2(b'*', b'[', path.as_bytes()).map(|i| {
26 path.as_bytes()[..i]
27 .iter()
28 .rposition(|x| *x == b'/')
29 .map_or(0, |x| 1 + x)
30 }) {
31 replacements.push((after_last_slash, 0, &[]));
33 after_last_slash
34 } else {
35 usize::MAX
36 };
37
38 while pos < path.len() {
39 match memchr::memchr2(b'*', b'.', &path.as_bytes()[pos..]) {
40 None => break,
41 Some(i) => pos += i,
42 }
43
44 let (len, replace): (usize, &[u8]) = match &path[pos..] {
45 v if v.starts_with("**") && (v.len() == 2 || v.as_bytes()[2] == b'/') => {
49 (3, b"(.*/)?" as _)
51 },
52 v if v.starts_with("**") => {
53 polars_bail!(ComputeError: "invalid ** glob pattern")
54 },
55 v if v.starts_with('*') => (1, b"[^/]*" as _),
56 v if v.starts_with('.') => (1, b"\\." as _),
58 _ => {
59 pos += 1;
60 continue;
61 },
62 };
63
64 replacements.push((pos, len, replace));
65 pos += len;
66 }
67
68 if replacements.is_empty() {
69 return Ok((Cow::Borrowed(path), None));
70 }
71
72 let prefix = Cow::Borrowed(&path[..replacements[0].0]);
73
74 let mut pos = replacements[0].0;
75 let mut expansion = Vec::with_capacity(path.len() - pos);
76 expansion.push(b'^');
77
78 for (offset, len, replace) in replacements {
79 expansion.extend_from_slice(&path.as_bytes()[pos..offset]);
80 expansion.extend_from_slice(replace);
81 pos = offset + len;
82 }
83
84 if pos < path.len() {
85 expansion.extend_from_slice(&path.as_bytes()[pos..]);
86 }
87
88 expansion.push(b'$');
89
90 Ok((prefix, Some(String::from_utf8(expansion).unwrap())))
91}
92
93#[derive(PartialEq, Debug, Default)]
95pub struct CloudLocation {
96 pub scheme: &'static str,
98 pub bucket: PlSmallStr,
100 pub prefix: String,
102 pub expansion: Option<PlSmallStr>,
104}
105
106impl CloudLocation {
107 pub fn new(path: PlRefPath, glob: bool) -> PolarsResult<Self> {
108 if let Some(scheme @ CloudScheme::Http | scheme @ CloudScheme::Https) = path.scheme() {
109 return Ok(CloudLocation {
111 scheme: scheme.as_str(),
112 ..Default::default()
113 });
114 }
115
116 let path_is_local = matches!(
117 path.scheme(),
118 None | Some(CloudScheme::File | CloudScheme::FileNoHostname)
119 );
120
121 let (bucket, key) = path
122 .strip_scheme_split_authority()
123 .ok_or(Cow::Borrowed(
124 "could not extract bucket/key (path did not contain '/')",
125 ))
126 .and_then(|x @ (bucket, _)| {
127 let bucket_is_empty = bucket.is_empty();
128
129 if path_is_local && !bucket_is_empty {
130 Err(Cow::Owned(format!(
131 "unsupported: non-empty hostname for 'file:' URI: '{bucket}'",
132 )))
133 } else if bucket_is_empty && !path_is_local {
134 Err(Cow::Borrowed("empty bucket name"))
135 } else {
136 Ok(x)
137 }
138 })
139 .map_err(|failed_reason| {
140 polars_err!(
141 ComputeError:
142 "failed to create CloudLocation: {} (path: '{}')",
143 failed_reason,
144 path,
145 )
146 })?;
147
148 let key = if path_is_local {
149 key
150 } else {
151 key.strip_prefix('/').unwrap_or(key)
152 };
153
154 let (prefix, expansion) = if glob {
155 let (prefix, expansion) = extract_prefix_expansion(key)?;
156
157 assert_eq!(prefix.starts_with('/'), key.starts_with('/'));
158
159 (prefix, expansion.map(|x| x.into()))
160 } else {
161 (key.into(), None)
162 };
163
164 Ok(CloudLocation {
165 scheme: path.scheme().unwrap_or(CloudScheme::File).as_str(),
166 bucket: PlSmallStr::from_str(bucket),
167 prefix: prefix.into_owned(),
168 expansion,
169 })
170 }
171}
172
173fn full_url(scheme: &str, bucket: &str, key: Path) -> String {
175 format!("{scheme}://{bucket}/{key}")
176}
177
178pub(crate) struct Matcher {
181 prefix: String,
182 re: Option<Regex>,
183}
184
185impl Matcher {
186 pub(crate) fn new(prefix: String, expansion: Option<&str>) -> PolarsResult<Matcher> {
188 let re = expansion
190 .map(polars_utils::regex_cache::compile_regex)
191 .transpose()?;
192 Ok(Matcher { prefix, re })
193 }
194
195 pub(crate) fn is_matching(&self, key: &str) -> bool {
196 if !key.starts_with(self.prefix.as_str()) {
197 return false;
199 }
200 if self.re.is_none() {
201 return true;
202 }
203 let last = &key[self.prefix.len()..];
204 self.re.as_ref().unwrap().is_match(last.as_ref())
205 }
206}
207
208pub async fn glob(
210 url: PlRefPath,
211 cloud_options: Option<&CloudOptions>,
212) -> PolarsResult<Vec<(String, u64)>> {
213 let (
216 CloudLocation {
217 scheme,
218 bucket,
219 prefix,
220 expansion,
221 },
222 store,
223 ) = super::build_object_store(url, cloud_options, true).await?;
224 let matcher = &Matcher::new(
225 if scheme == "file" {
226 prefix[1..].into()
228 } else {
229 prefix.clone()
230 },
231 expansion.as_deref(),
232 )?;
233
234 let path = Path::from(prefix.as_str());
235 let path = Some(&path);
236
237 let mut locations = store
238 .exec_with_rebuild_retry_on_err(|store| async move {
239 store
240 .list(path)
241 .try_filter_map(|x| async move {
242 let out = (x.size > 0 && matcher.is_matching(x.location.as_ref()))
244 .then_some((x.location, x.size));
245 Ok(out)
246 })
247 .try_collect::<Vec<_>>()
248 .await
249 })
250 .await?;
251
252 locations.sort_unstable_by(|a, b| a.0.cmp(&b.0));
253 Ok(locations
254 .into_iter()
255 .map(|(l, size)| (full_url(scheme, &bucket, l), size))
256 .collect::<Vec<_>>())
257}
258
259#[cfg(test)]
260mod test {
261 use super::*;
262
263 #[test]
264 fn test_cloud_location() {
265 assert_eq!(
266 CloudLocation::new(PlRefPath::new("s3://a/b"), true).unwrap(),
267 CloudLocation {
268 scheme: "s3",
269 bucket: "a".into(),
270 prefix: "b".into(),
271 expansion: None,
272 }
273 );
274 assert_eq!(
275 CloudLocation::new(PlRefPath::new("s3://a/b/*.c"), true).unwrap(),
276 CloudLocation {
277 scheme: "s3",
278 bucket: "a".into(),
279 prefix: "b/".into(),
280 expansion: Some("^[^/]*\\.c$".into()),
281 }
282 );
283 assert_eq!(
284 CloudLocation::new(PlRefPath::new("file:///a/b"), true).unwrap(),
285 CloudLocation {
286 scheme: "file",
287 bucket: "".into(),
288 prefix: "/a/b".into(),
289 expansion: None,
290 }
291 );
292 assert_eq!(
293 CloudLocation::new(PlRefPath::new("file:/a/b"), true).unwrap(),
294 CloudLocation {
295 scheme: "file",
296 bucket: "".into(),
297 prefix: "/a/b".into(),
298 expansion: None,
299 }
300 );
301 }
302
303 #[test]
304 fn test_extract_prefix_expansion() {
305 assert!(extract_prefix_expansion("**url").is_err());
306 assert_eq!(
307 extract_prefix_expansion("a/b.c").unwrap(),
308 ("a/b.c".into(), None)
309 );
310 assert_eq!(
311 extract_prefix_expansion("a/**").unwrap(),
312 ("a/".into(), Some("^(.*/)?$".into()))
313 );
314 assert_eq!(
315 extract_prefix_expansion("a/**/b").unwrap(),
316 ("a/".into(), Some("^(.*/)?b$".into()))
317 );
318 assert_eq!(
319 extract_prefix_expansion("a/**/*b").unwrap(),
320 ("a/".into(), Some("^(.*/)?[^/]*b$".into()))
321 );
322 assert_eq!(
323 extract_prefix_expansion("a/**/data/*b").unwrap(),
324 ("a/".into(), Some("^(.*/)?data/[^/]*b$".into()))
325 );
326 assert_eq!(
327 extract_prefix_expansion("a/*b").unwrap(),
328 ("a/".into(), Some("^[^/]*b$".into()))
329 );
330 }
331
332 #[test]
333 fn test_matcher_file_name() {
334 let cloud_location =
335 CloudLocation::new(PlRefPath::new("s3://bucket/folder/*.parquet"), true).unwrap();
336 let a = Matcher::new(cloud_location.prefix, cloud_location.expansion.as_deref()).unwrap();
337 assert!(a.is_matching(Path::from("folder/1.parquet").as_ref()));
339 assert!(!a.is_matching(Path::from("folder/1parquet").as_ref()));
341 assert!(!a.is_matching(Path::from("folder/other/1.parquet").as_ref()));
343 }
344
345 #[test]
346 fn test_matcher_folders() {
347 let cloud_location =
348 CloudLocation::new(PlRefPath::new("s3://bucket/folder/**/*.parquet"), true).unwrap();
349
350 let a = Matcher::new(cloud_location.prefix, cloud_location.expansion.as_deref()).unwrap();
351 assert!(a.is_matching(Path::from("folder/1.parquet").as_ref()));
353 assert!(a.is_matching(Path::from("folder/other/1.parquet").as_ref()));
355
356 let cloud_location =
357 CloudLocation::new(PlRefPath::new("s3://bucket/folder/**/data/*.parquet"), true)
358 .unwrap();
359 let a = Matcher::new(cloud_location.prefix, cloud_location.expansion.as_deref()).unwrap();
360
361 assert!(!a.is_matching(Path::from("folder/1.parquet").as_ref()));
363 assert!(a.is_matching(Path::from("folder/data/1.parquet").as_ref()));
365 assert!(a.is_matching(Path::from("folder/other/data/1.parquet").as_ref()));
367 }
368
369 #[test]
370 fn test_cloud_location_no_glob() {
371 let cloud_location = CloudLocation::new(PlRefPath::new("s3://bucket/[*"), false).unwrap();
372 assert_eq!(
373 cloud_location,
374 CloudLocation {
375 scheme: "s3",
376 bucket: "bucket".into(),
377 prefix: "[*".into(),
378 expansion: None,
379 },
380 )
381 }
382
383 #[test]
384 fn test_cloud_location_percentages() {
385 use super::CloudLocation;
386
387 let path = "s3://bucket/%25";
388 let cloud_location = CloudLocation::new(PlRefPath::new(path), true).unwrap();
389
390 assert_eq!(
391 cloud_location,
392 CloudLocation {
393 scheme: "s3",
394 bucket: "bucket".into(),
395 prefix: "%25".into(),
396 expansion: None,
397 }
398 );
399
400 let path = "https://pola.rs/%25";
401 let cloud_location = CloudLocation::new(PlRefPath::new(path), true).unwrap();
402
403 assert_eq!(
404 cloud_location,
405 CloudLocation {
406 scheme: "https",
407 bucket: "".into(),
408 prefix: "".into(),
409 expansion: None,
410 }
411 );
412 }
413
414 #[test]
415 fn test_glob_wildcard_21736() {
416 let path = "s3://bucket/folder/**/data.parquet";
417 let cloud_location = CloudLocation::new(PlRefPath::new(path), true).unwrap();
418
419 let a = Matcher::new(cloud_location.prefix, cloud_location.expansion.as_deref()).unwrap();
420
421 assert!(!a.is_matching("folder/_data.parquet"));
422
423 assert!(a.is_matching("folder/data.parquet"));
424 assert!(a.is_matching("folder/abc/data.parquet"));
425 assert!(a.is_matching("folder/abc/def/data.parquet"));
426
427 let path = "s3://bucket/folder/data_*.parquet";
428 let cloud_location = CloudLocation::new(PlRefPath::new(path), true).unwrap();
429
430 let a = Matcher::new(cloud_location.prefix, cloud_location.expansion.as_deref()).unwrap();
431
432 assert!(!a.is_matching("folder/data_1.ipc"))
433 }
434}