1use std::path::{Path, PathBuf};
8
9use polars::prelude::*;
10
11use crate::app::modals::filter_modal::{FilterOperator, FilterStatement, LogicalOperator};
12use crate::app::modals::pivot_melt_modal::PivotAggregation;
13use crate::{CompressionFormat, FileFormat, OpenOptions};
14
15pub(crate) fn py_str(s: &str) -> String {
17 let mut out = String::with_capacity(s.len() + 2);
18 out.push('"');
19 for c in s.chars() {
20 match c {
21 '\\' => out.push_str("\\\\"),
22 '"' => out.push_str("\\\""),
23 '\n' => out.push_str("\\n"),
24 '\r' => out.push_str("\\r"),
25 '\t' => out.push_str("\\t"),
26 c if c.is_control() => out.push_str(&format!("\\u{:04x}", c as u32)),
27 c => out.push(c),
28 }
29 }
30 out.push('"');
31 out
32}
33
34pub(crate) fn py_comment(text: &str) -> String {
37 let mut out = String::with_capacity(text.len() + 2);
38 out.push_str("# ");
39 for c in text.chars() {
40 match c {
41 '\n' => out.push_str("\\n"),
42 '\r' => out.push_str("\\r"),
43 '\t' => out.push('\t'),
44 c if c.is_control() || c == '\u{2028}' || c == '\u{2029}' => {
45 out.push_str(&format!("\\u{:04x}", c as u32))
46 }
47 c => out.push(c),
48 }
49 }
50 out
51}
52
53pub(crate) fn py_float(f: f64) -> String {
56 if f.is_nan() {
57 "float(\"nan\")".to_string()
58 } else if f.is_infinite() {
59 if f > 0.0 {
60 "float(\"inf\")".to_string()
61 } else {
62 "float(\"-inf\")".to_string()
63 }
64 } else {
65 format!("{f:?}")
67 }
68}
69
70pub(crate) fn py_bool(b: bool) -> &'static str {
71 if b { "True" } else { "False" }
72}
73
74pub(crate) fn py_names(names: &[String]) -> String {
76 let items: Vec<String> = names.iter().map(|n| py_str(n)).collect();
77 format!("[{}]", items.join(", "))
78}
79
80pub(crate) fn sort_call(columns: &[String], descending: &[bool]) -> String {
82 let by = match columns {
83 [one] => py_str(one),
84 _ => py_names(columns),
85 };
86 let descending = if descending.iter().all(|d| !d) {
87 String::new()
88 } else if descending.iter().all(|d| *d) {
89 "descending=True, ".to_string()
90 } else {
91 let flags: Vec<&str> = descending.iter().map(|d| py_bool(*d)).collect();
92 format!("descending=[{}], ", flags.join(", "))
93 };
94 format!(".sort({by}, {descending}nulls_last=True, maintain_order=True)")
95}
96
97#[derive(Debug, Clone, PartialEq)]
101pub enum FilterValue {
102 Typed(Scalar),
103 Str(String),
104}
105
106impl FilterValue {
107 fn lit(&self) -> Expr {
108 match self {
109 FilterValue::Typed(scalar) => lit(scalar.clone()),
110 FilterValue::Str(s) => lit(s.as_str()),
111 }
112 }
113
114 fn python(&self) -> String {
115 match self {
116 FilterValue::Typed(scalar) => crate::typed_value::python(scalar),
117 FilterValue::Str(s) => py_str(s),
118 }
119 }
120}
121
122#[derive(Debug, Clone, PartialEq)]
125pub struct SidebarFilter {
126 pub column: String,
127 pub operator: FilterOperator,
128 pub value: FilterValue,
129 pub text: String,
131 pub logical_op: LogicalOperator,
132 pub searched: Vec<(String, DataType)>,
135}
136
137impl SidebarFilter {
138 pub fn typed_in(statement: &FilterStatement, schema: &Schema, shown: &[String]) -> Self {
142 let mut filter = Self::typed(statement, schema.get(&statement.column));
143 if statement.operator.is_find() {
144 let spec = filter.find_spec();
145 let names: Vec<&String> =
146 if statement.column == crate::app::modals::filter_modal::ANY_COLUMN {
147 if statement.columns.is_empty() {
148 shown.iter().collect()
149 } else {
150 statement.columns.iter().collect()
151 }
152 } else {
153 vec![&statement.column]
154 };
155 filter.searched = names
156 .into_iter()
157 .filter_map(|name| Some((name.clone(), schema.get(name)?.clone())))
158 .filter(|(name, dtype)| crate::find::cell_matches(&spec, name, dtype).is_some())
159 .collect();
160 }
161 filter
162 }
163
164 pub fn unscriptable_columns(&self) -> Vec<String> {
167 if !self.operator.is_find() {
168 return Vec::new();
169 }
170 self.searched
171 .iter()
172 .filter(|(_, dtype)| matches!(dtype, DataType::Duration(_)))
173 .map(|(name, _)| name.clone())
174 .collect()
175 }
176
177 fn find_spec(&self) -> crate::find::FindSpec {
179 crate::find::FindSpec {
180 pattern: self.text.clone(),
181 regex: self.operator == FilterOperator::HasRegex,
182 fuzzy: self.operator == FilterOperator::HasFuzzy,
183 column: None,
184 }
185 }
186
187 pub fn typed(statement: &FilterStatement, dtype: Option<&DataType>) -> Self {
188 let text = statement.value.as_str();
189 let value = match dtype {
190 None | Some(DataType::String) => FilterValue::Str(text.to_string()),
191 Some(dtype) => crate::typed_value::parse(text, dtype)
192 .map(FilterValue::Typed)
193 .unwrap_or_else(|_| FilterValue::Str(text.to_string())),
194 };
195 Self {
196 column: statement.column.clone(),
197 operator: statement.operator,
198 value,
199 text: statement.value.clone(),
200 logical_op: statement.logical_op,
201 searched: Vec::new(),
202 }
203 }
204
205 pub fn problem(statement: &FilterStatement, dtype: Option<&DataType>) -> Option<String> {
209 if statement.operator == FilterOperator::HasRegex
211 && let Err(reason) = Self::typed(statement, dtype).find_spec().check()
212 {
213 return Some(reason);
214 }
215 if statement.operator.is_find()
218 && let Some(dtype) = dtype
219 && crate::find::cell_matches(
220 &Self::typed(statement, Some(dtype)).find_spec(),
221 &statement.column,
222 dtype,
223 )
224 .is_none()
225 {
226 return Some(format!("{}: no text to match", statement.column));
227 }
228 let compares = statement.operator.takes_value()
229 && !statement.operator.is_find()
230 && !matches!(
231 statement.operator,
232 FilterOperator::Contains | FilterOperator::NotContains
233 );
234 let dtype = dtype.filter(|_| compares)?;
235 crate::typed_value::parse(&statement.value, dtype)
236 .err()
237 .map(|why| format!("{}: {why}", statement.column))
238 }
239
240 fn expr(&self) -> Expr {
241 let column = col(&self.column);
242 let contains = || {
243 col(&self.column)
244 .str()
245 .contains_literal(lit(self.text.as_str()))
246 };
247 match self.operator {
248 FilterOperator::Eq => column.eq(self.value.lit()),
249 FilterOperator::NotEq => column.neq(self.value.lit()),
250 FilterOperator::Gt => column.gt(self.value.lit()),
251 FilterOperator::Lt => column.lt(self.value.lit()),
252 FilterOperator::GtEq => column.gt_eq(self.value.lit()),
253 FilterOperator::LtEq => column.lt_eq(self.value.lit()),
254 FilterOperator::Contains => contains(),
255 FilterOperator::NotContains => contains().not(),
256 FilterOperator::IsNull => column.is_null(),
257 FilterOperator::IsNotNull => column.is_not_null(),
258 FilterOperator::Has | FilterOperator::HasRegex | FilterOperator::HasFuzzy => {
259 let spec = self.find_spec();
260 let cells: Vec<Expr> = self
261 .searched
262 .iter()
263 .filter_map(|(name, dtype)| crate::find::cell_matches(&spec, name, dtype))
264 .collect();
265 cells
266 .into_iter()
267 .reduce(|a, b| a.or(b))
268 .unwrap_or(lit(false))
269 }
270 }
271 }
272
273 fn python(&self) -> String {
274 let column = format!("pl.col({})", py_str(&self.column));
275 let op = match self.operator {
276 FilterOperator::Eq => "==",
277 FilterOperator::NotEq => "!=",
278 FilterOperator::Gt => ">",
279 FilterOperator::Lt => "<",
280 FilterOperator::GtEq => ">=",
281 FilterOperator::LtEq => "<=",
282 FilterOperator::Contains | FilterOperator::NotContains => {
283 let not = if self.operator == FilterOperator::NotContains {
284 "~"
285 } else {
286 ""
287 };
288 return format!(
289 "{not}{column}.str.contains({}, literal=True)",
290 py_str(&self.text)
291 );
292 }
293 FilterOperator::IsNull => return format!("{column}.is_null()"),
294 FilterOperator::IsNotNull => return format!("{column}.is_not_null()"),
295 FilterOperator::Has | FilterOperator::HasRegex | FilterOperator::HasFuzzy => {
296 let spec = self.find_spec();
297 let cells: Vec<String> = self
300 .searched
301 .iter()
302 .filter(|(_, dtype)| !matches!(dtype, DataType::Duration(_)))
303 .map(|(name, _)| {
304 let column = format!("pl.col({}).cast(pl.String)", py_str(name));
305 match spec.regex_source() {
306 None => format!(
307 "{column}.str.contains({}, literal=True)",
308 py_str(&spec.pattern)
309 ),
310 Some(source) => {
311 format!("{column}.str.contains({})", py_str(&source))
312 }
313 }
314 })
315 .collect();
316 return match cells.len() {
317 0 => "pl.lit(False)".to_string(),
318 1 => format!("{}.fill_null(False)", cells[0]),
319 _ => format!("pl.any_horizontal({}).fill_null(False)", cells.join(", ")),
320 };
321 }
322 };
323 format!("{column} {op} {}", self.value.python())
324 }
325}
326
327pub fn filters_expr(filters: &[SidebarFilter]) -> Option<Expr> {
330 filters.iter().fold(None, |all, f| {
331 Some(match all {
332 None => f.expr(),
333 Some(all) => match f.logical_op {
334 LogicalOperator::And => all.and(f.expr()),
335 LogicalOperator::Or => all.or(f.expr()),
336 },
337 })
338 })
339}
340
341fn filters_python(filters: &[SidebarFilter]) -> String {
342 let mut out = String::new();
343 let mut last: Option<LogicalOperator> = None;
344 for (i, f) in filters.iter().enumerate() {
345 let term = format!("({})", f.python());
346 if i == 0 {
347 out = if filters.len() == 1 { f.python() } else { term };
349 continue;
350 }
351 if last.is_some_and(|l| l != f.logical_op) {
354 out = format!("({out})");
355 }
356 let op = match f.logical_op {
357 LogicalOperator::And => "&",
358 LogicalOperator::Or => "|",
359 };
360 out = format!("{out} {op} {term}");
361 last = Some(f.logical_op);
362 }
363 out
364}
365
366#[derive(Debug, Clone, PartialEq)]
368pub enum Step {
369 Query {
372 query: String,
373 input: SchemaRef,
374 keys: Vec<String>,
375 },
376 QueryRows {
378 query: String,
379 input: SchemaRef,
380 },
381 Sql {
384 sql: String,
385 ordered_by: Vec<String>,
386 },
387 Search {
389 patterns: Vec<String>,
390 columns: Vec<String>,
391 },
392 Filter(Vec<SidebarFilter>),
393 Sort {
394 columns: Vec<String>,
395 descending: Vec<bool>,
396 },
397 Reverse,
398 Select(Vec<String>),
399 Drop(Vec<String>),
400 Pivot {
401 index: Vec<String>,
402 on: String,
403 values: String,
404 aggregation: PivotAggregation,
405 },
406 Melt {
407 index: Vec<String>,
408 on: Vec<String>,
409 variable_name: String,
410 value_name: String,
411 },
412 Matching(Vec<(String, String)>),
415 Unreproducible(String),
417}
418
419impl Step {
420 fn python(&self) -> Vec<String> {
422 match self {
423 Step::Query { query, input, keys } => match crate::query::parse_nodes(query) {
424 Ok(mut nodes) => {
425 nodes.resolve_division(input);
426 nodes.resolve_time_zones(input);
427 nodes.python_steps(keys)
428 }
429 Err(e) => vec![py_comment(&format!("the query did not parse: {e}"))],
430 },
431 Step::QueryRows { query, input } => match crate::query::parse_nodes(query) {
432 Ok(mut nodes) => {
433 nodes.resolve_division(input);
434 nodes.resolve_time_zones(input);
435 nodes.python_filter().into_iter().collect()
436 }
437 Err(e) => vec![py_comment(&format!("the query did not parse: {e}"))],
438 },
439 Step::Sql { sql, ordered_by } => {
440 let sql = sql.trim();
441 let verbatim = sql.contains('\n')
444 && !sql.contains("\"\"\"")
445 && !sql.ends_with('"')
446 && sql
447 .chars()
448 .all(|c| c == '\n' || c == '\t' || (c != '\\' && !c.is_control()));
449 let mut lines = if verbatim {
450 vec![
451 ".sql(".to_string(),
452 format!(" \"\"\"{sql}\"\"\","),
453 " table_name=\"df\",".to_string(),
454 ")".to_string(),
455 ]
456 } else {
457 vec![format!(".sql({}, table_name=\"df\")", py_str(sql))]
458 };
459 if !ordered_by.is_empty() {
460 lines.push(sort_call(ordered_by, &vec![false; ordered_by.len()]));
461 }
462 lines
463 }
464 Step::Search { patterns, columns } => {
465 let terms: Vec<String> = patterns
466 .iter()
467 .map(|p| {
468 let any: Vec<String> = columns
469 .iter()
470 .map(|c| {
471 format!(
472 "pl.col({}).str.contains({}, strict=False)",
473 py_str(c),
474 py_str(p)
475 )
476 })
477 .collect();
478 if any.len() == 1 || patterns.len() == 1 {
479 any.join(" | ")
480 } else {
481 format!("({})", any.join(" | "))
482 }
483 })
484 .collect();
485 vec![format!(".filter({})", terms.join(" & "))]
486 }
487 Step::Filter(filters) => vec![format!(".filter({})", filters_python(filters))],
488 Step::Sort {
489 columns,
490 descending,
491 } => vec![sort_call(columns, descending)],
492 Step::Reverse => vec![".reverse()".to_string()],
493 Step::Select(columns) => vec![format!(".select({})", py_names(columns))],
494 Step::Drop(columns) => vec![format!(".drop({})", py_names(columns))],
495 Step::Pivot {
496 index,
497 on,
498 values,
499 aggregation,
500 } => {
501 let agg = match aggregation {
504 PivotAggregation::Last => "last()",
505 PivotAggregation::First => "first()",
506 PivotAggregation::Min => "min()",
507 PivotAggregation::Max => "max()",
508 PivotAggregation::Avg => "mean()",
509 PivotAggregation::Med => "median()",
510 PivotAggregation::Std => "std()",
511 PivotAggregation::Count => "len()",
512 };
513 let cell = match aggregation {
514 PivotAggregation::Count => "sum",
515 _ => "first",
516 };
517 let keys: Vec<String> = index.iter().chain([on]).cloned().collect();
518 vec![
519 format!(".group_by({}, maintain_order=True)", py_names(&keys)),
520 format!(".agg(pl.col({}).{agg})", py_str(values)),
521 ".collect()".to_string(),
522 ".pipe(".to_string(),
523 " lambda cells: cells.pivot(".to_string(),
524 format!(" on={},", py_str(on)),
525 format!(
526 " on_columns=cells[{}].unique().sort(nulls_last=True),",
527 py_str(on)
528 ),
529 format!(" index={},", py_names(index)),
530 format!(" values={},", py_str(values)),
531 format!(" aggregate_function={},", py_str(cell)),
532 " )".to_string(),
533 ")".to_string(),
534 ".lazy()".to_string(),
535 ]
536 }
537 Step::Melt {
538 index,
539 on,
540 variable_name,
541 value_name,
542 } => vec![format!(
543 ".unpivot(on={}, index={}, variable_name={}, value_name={})",
544 py_names(on),
545 py_names(index),
546 py_str(variable_name),
547 py_str(value_name)
548 )],
549 Step::Matching(keys) => {
550 let terms: Vec<String> = keys
551 .iter()
552 .map(|(key, value)| format!("{key}.eq_missing({value})"))
553 .collect();
554 let terms = if terms.len() == 1 {
555 terms
556 } else {
557 terms.into_iter().map(|t| format!("({t})")).collect()
558 };
559 vec![format!(".filter({})", terms.join(" & "))]
560 }
561 Step::Unreproducible(what) => vec![py_comment(what)],
562 }
563 }
564}
565
566#[derive(Debug, Clone, PartialEq)]
568pub enum Source {
569 Read {
573 call: String,
574 after: Vec<String>,
575 notes: Vec<String>,
576 imports: Vec<&'static str>,
578 },
579 Placeholder { what: String },
582}
583
584pub struct OpenRecord<'a> {
586 pub paths: Option<&'a [PathBuf]>,
588 pub options: &'a OpenOptions,
589 pub format: Option<FileFormat>,
593 pub read_mode: Option<crate::ReadMode>,
595 pub schema: &'a Schema,
597 pub remote_objects: Vec<String>,
599 pub s3_endpoint: Option<String>,
602 pub s3_region: Option<String>,
603 pub unsigned: bool,
605 pub read_as_text: Vec<String>,
607 pub spec: Option<String>,
609}
610
611fn is_url(path: &Path) -> bool {
612 crate::cloud::source::is_remote_url(path)
613}
614
615fn without_secrets(url: &str) -> (String, bool) {
618 let url = &*crate::cloud::source::split_source_id(url).1;
620 let Some(scheme_end) = url.find("://").map(|i| i + 3) else {
621 return (url.to_string(), false);
622 };
623 let (scheme, rest) = url.split_at(scheme_end);
624 let host_end = rest.find(['/', '?', '#']).unwrap_or(rest.len());
625 let (authority, path) = rest.split_at(host_end);
626 let azure = ["abfs://", "abfss://"]
628 .iter()
629 .any(|s| scheme.eq_ignore_ascii_case(s));
630 let host = match authority.rsplit_once('@') {
631 Some((_, host)) if !azure => host,
632 _ => authority,
633 };
634 let http = scheme.eq_ignore_ascii_case("http://") || scheme.eq_ignore_ascii_case("https://");
635 let path = match path.find(['?', '#']) {
636 Some(i) if http => &path[..i],
637 _ => path,
638 };
639 let kept = format!("{scheme}{host}{path}");
640 let cut = kept != url;
641 (kept, cut)
642}
643
644fn file_format(path: &Path, record: &OpenRecord) -> Option<FileFormat> {
647 record.format.or(record.options.format).or_else(|| {
648 FileFormat::from_path(path).or_else(|| {
649 CompressionFormat::from_extension(path)
650 .and_then(|_| path.file_stem())
651 .and_then(|stem| FileFormat::from_path(Path::new(stem)))
652 })
653 })
654}
655
656fn commonest_format<'a>(names: impl Iterator<Item = &'a str>) -> Option<FileFormat> {
658 let mut counts: Vec<(FileFormat, usize)> = Vec::new();
659 for name in names {
660 if let Some(format) = FileFormat::from_path(Path::new(name)) {
661 match counts.iter_mut().find(|(f, _)| *f == format) {
662 Some((_, n)) => *n += 1,
663 None => counts.push((format, 1)),
664 }
665 }
666 }
667 counts.into_iter().max_by_key(|(_, n)| *n).map(|(f, _)| f)
668}
669
670struct Target {
672 text: String,
675 format: FileFormat,
676 below: bool,
678 pattern: bool,
680 literal: bool,
682}
683
684fn reader_target(path: &Path, record: &OpenRecord) -> Option<Target> {
686 let text = path.to_string_lossy().to_string();
687 if is_url(path) {
688 let (text, _) = without_secrets(&text);
689 if let Some(format) = file_format(Path::new(&text), record) {
690 let pattern = crate::cloud::source::has_glob_chars(Path::new(&text));
691 return Some(Target {
692 text,
693 format,
694 below: false,
695 pattern,
696 literal: false,
697 });
698 }
699 let format = record
701 .format
702 .or(record.options.format)
703 .or_else(|| commonest_format(record.remote_objects.iter().map(String::as_str)))?;
704 let base = text.trim_end_matches('/');
705 let ext = format_extension(format)?;
706 return Some(Target {
707 text: format!("{base}/**/*.{ext}"),
708 format,
709 below: true,
710 pattern: true,
711 literal: false,
712 });
713 }
714 if path.is_dir() {
715 let entries: Vec<std::fs::DirEntry> = std::fs::read_dir(path).ok()?.flatten().collect();
716 let mut names: Vec<String> = entries
717 .iter()
718 .filter(|e| e.path().is_file())
719 .map(|e| e.file_name().to_string_lossy().to_string())
720 .collect();
721 if names
723 .iter()
724 .any(|n| FileFormat::from_path(Path::new(n)) == Some(FileFormat::Arrow))
725 {
726 names.retain(|n| !crate::home::discover::is_hugging_face_metadata(n));
727 }
728 let has_dirs = entries.iter().any(|e| e.path().is_dir());
729 let format = record
730 .format
731 .or(record.options.format)
732 .or_else(|| commonest_format(names.iter().map(String::as_str)));
733 let base = crate::cloud::source::escape_glob(text.trim_end_matches(['/', '\\']));
735 let (text, format, below) = match format {
738 Some(FileFormat::Parquet) | None if has_dirs => {
739 (format!("{base}/**/*.parquet"), FileFormat::Parquet, true)
740 }
741 Some(format) => (
742 format!("{base}/*.{}", format_extension(format)?),
743 format,
744 false,
745 ),
746 None => return None,
747 };
748 return Some(Target {
749 text,
750 format,
751 below,
752 pattern: true,
753 literal: false,
754 });
755 }
756 let format = file_format(path, record)?;
757 let pattern = crate::cloud::source::expands_as_glob(path);
758 Some(Target {
759 literal: !pattern && scans_by_pattern(format) && crate::cloud::source::has_glob_chars(path),
762 text,
763 format,
764 below: false,
765 pattern,
766 })
767}
768
769fn scans_by_pattern(format: FileFormat) -> bool {
771 python_of(format).is_some_and(|python| !python.eager)
772}
773
774fn hugging_face_files(paths: &[PathBuf], table: Option<&str>) -> Option<Vec<String>> {
777 let [path] = paths else {
778 return None;
779 };
780 if is_url(path) || !path.is_dir() {
781 return None;
782 }
783 let dict_split = crate::formats::hf_splits::dataset_dict(path).and_then(|splits| {
785 let listed: Vec<&str> = splits.iter().map(String::as_str).collect();
786 crate::formats::hf_splits::pick(&listed, table).ok()?.split
787 });
788 let dict = dict_split.is_some();
789 let table = if dict { None } else { table };
790 let path = &dict_split.map_or_else(|| path.clone(), |split| path.join(split));
791 let mut inside: Vec<PathBuf> = std::fs::read_dir(path)
792 .ok()?
793 .flatten()
794 .map(|e| e.path())
795 .filter(|p| p.is_file() && FileFormat::from_path(p) == Some(FileFormat::Arrow))
796 .collect();
797 inside.sort();
798 let cache = ["dataset_info.json", "state.json"]
799 .iter()
800 .any(|name| path.join(name).is_file());
801 if !cache && !dict {
802 return None;
803 }
804 if cache {
805 let names: Vec<&str> = inside
806 .iter()
807 .map(|f| f.file_name().and_then(|n| n.to_str()).unwrap_or_default())
808 .collect();
809 let (chosen, _) = crate::formats::hf_splits::choose(&names, table).ok()?;
810 inside = chosen.into_iter().map(|i| inside[i].clone()).collect();
811 }
812 Some(
813 inside
814 .iter()
815 .map(|p| p.to_string_lossy().to_string())
816 .collect(),
817 )
818}
819
820fn arrow_read(inputs: &[(String, bool)], extra: Option<&str>) -> String {
824 let extra = extra.map(|e| format!(", {e}")).unwrap_or_default();
825 let names: Vec<String> = inputs.iter().map(|(name, _)| name.clone()).collect();
826 if inputs.iter().all(|(_, stream)| *stream) {
827 return match names.as_slice() {
828 [one] => format!("pl.read_ipc_stream({}{extra}).lazy()", py_str(one)),
829 many => format!(
830 "pl.concat([pl.read_ipc_stream(f{extra}) for f in {}]).lazy()",
831 py_names(many)
832 ),
833 };
834 }
835 if inputs.iter().all(|(_, stream)| !*stream) {
836 return match names.as_slice() {
837 [one] => format!("pl.scan_ipc({}{extra})", py_str(one)),
838 many => format!("pl.scan_ipc({}{extra})", py_names(many)),
839 };
840 }
841 let reads: Vec<String> = inputs
842 .iter()
843 .map(|(name, stream)| match stream {
844 true => format!("pl.read_ipc_stream({}{extra}).lazy()", py_str(name)),
845 false => format!("pl.scan_ipc({}{extra})", py_str(name)),
846 })
847 .collect();
848 format!(
849 "pl.concat([{}], how=\"diagonal_relaxed\")",
850 reads.join(", ")
851 )
852}
853
854fn format_extension(format: FileFormat) -> Option<&'static str> {
857 python_of(format)?;
858 let d = format.descriptor();
859 (d.many_files || format.separator().is_some())
860 .then(|| d.extensions.first().copied())
861 .flatten()
862}
863
864pub(crate) struct Python {
867 pub call: &'static str,
869 pub eager: bool,
871 pub glob_flag: bool,
873 pub arguments: Option<fn(&mut Call<'_>) -> Option<Source>>,
876}
877
878pub(crate) struct Call<'a> {
880 pub record: &'a OpenRecord<'a>,
881 pub paths: &'a [PathBuf],
882 pub format: FileFormat,
883 pub names: &'a [String],
885 pub below: bool,
887 pub storage: Option<String>,
889 pub args: Vec<String>,
890 pub after: Vec<String>,
891 pub skip_tail: Option<String>,
893 pub notes: Vec<String>,
894}
895
896impl Call<'_> {
897 fn storage_for(&self, names: impl IntoIterator<Item = impl AsRef<str>>) -> Option<String> {
899 names
900 .into_iter()
901 .any(|n| store_scheme(n.as_ref()).is_some())
902 .then(|| self.storage.clone())
903 .flatten()
904 }
905}
906
907fn store_scheme(name: &str) -> Option<&'static str> {
910 let (scheme, _) = name.split_once("://")?;
911 match scheme.to_ascii_lowercase().as_str() {
912 "s3" | "s3a" => Some("s3"),
913 "gs" | "gcs" => Some("gs"),
914 "az" | "adl" | "azure" | "abfs" | "abfss" => Some("azure"),
915 _ => None,
916 }
917}
918
919fn storage_options(name: &str, record: &OpenRecord, endpoint: Option<&str>) -> Option<String> {
923 let mut pairs: Vec<(&str, String)> = Vec::new();
924 match store_scheme(name)? {
925 "s3" => {
926 pairs.extend(endpoint.map(|e| ("aws_endpoint_url", e.to_string())));
927 pairs.extend(record.s3_region.clone().map(|r| ("aws_region", r)));
928 }
929 "azure" => {
930 pairs.extend(
931 crate::cloud::source::azure_parts(name)
932 .map(|(account, ..)| ("account_name", account)),
933 );
934 }
935 _ => {}
936 }
937 if record.unsigned {
938 pairs.push(("skip_signature", "true".to_string()));
939 }
940 let pairs: Vec<String> = pairs
941 .iter()
942 .map(|(k, v)| format!("{}: {}", py_str(k), py_str(v)))
943 .collect();
944 (!pairs.is_empty()).then(|| format!("storage_options={{{}}}", pairs.join(", ")))
945}
946
947pub(crate) fn ndjson_arguments(call: &mut Call<'_>) -> Option<Source> {
949 if let Some(s) = call.storage_for(call.names) {
950 call.args.push(s);
951 }
952 None
953}
954
955fn python_of(format: FileFormat) -> Option<&'static Python> {
956 crate::formats::readers::of(format).python.as_ref()
957}
958
959pub(crate) fn parquet_arguments(call: &mut Call<'_>) -> Option<Source> {
961 if call.record.options.hive || call.below {
962 call.args.push("hive_partitioning=True".to_string());
963 }
964 if let Some(s) = call.storage_for(call.names) {
965 call.args.push(s);
966 }
967 None
968}
969
970pub(crate) fn csv_arguments(call: &mut Call<'_>) -> Option<Source> {
972 let options = call.record.options;
973 let names = call.names.join(", ");
974 let comment = options.comment_char.as_deref().filter(|c| !c.is_empty());
976 let dialect: Vec<&str> = [
977 (comment.is_some_and(|c| c.len() > 5), "--comment"),
978 (options.header_rows().is_some(), "--header-rows"),
979 (options.skip_initial_space, "--skip-initial-space"),
980 ]
981 .into_iter()
982 .filter_map(|(set, flag)| set.then_some(flag))
983 .collect();
984 if !dialect.is_empty() {
985 return Some(Source::Placeholder {
986 what: format!(
987 "{names}: datui read it with {}, which it cannot write as Python: load it here.",
988 dialect.join(", ")
989 ),
990 });
991 }
992 let separator = options
993 .delimiter
994 .or_else(|| call.format.separator())
995 .unwrap_or(b',');
996 let args = &mut call.args;
997 if separator != b',' {
998 args.push(format!(
999 "separator={}",
1000 py_str(&(separator as char).to_string())
1001 ));
1002 }
1003 if let Some(prefix) = comment {
1004 args.push(format!("comment_prefix={}", py_str(prefix)));
1005 }
1006 if options.has_header == Some(false) {
1007 args.push("has_header=False".to_string());
1008 }
1009 if let Some(n) = options.skip_lines {
1010 args.push(format!("skip_lines={n}"));
1011 }
1012 if let Some(n) = options.skip_rows {
1013 args.push(format!("skip_rows={n}"));
1014 }
1015 if let Some(n) = options.infer_schema_length {
1016 args.push(format!("infer_schema_length={n}"));
1017 }
1018 if options.ignore_errors {
1019 args.push("ignore_errors=True".to_string());
1020 }
1021 let bucket_prefix = call.below && call.paths.iter().any(|p| is_url(p));
1023 if options.csv_try_parse_dates() && !bucket_prefix {
1024 call.args.push("try_parse_dates=True".to_string());
1025 }
1026 if let Some(nulls) = csv_null_values(options, call.record.schema) {
1027 call.args.push(format!("null_values={nulls}"));
1028 }
1029 if !options.typing.text.is_empty() {
1031 let text: Vec<String> = options
1032 .typing
1033 .text
1034 .iter()
1035 .map(|name| format!("{}: pl.String", py_str(name)))
1036 .collect();
1037 call.args
1038 .push(format!("schema_overrides={{{}}}", text.join(", ")));
1039 }
1040 if let Some(s) = call.storage_for(call.names) {
1041 call.args.push(s);
1042 }
1043 if let Some(n) = options.skip_tail_rows.filter(|n| *n > 0) {
1044 call.skip_tail = Some(format!(".filter(pl.int_range(pl.len()) < pl.len() - {n})"));
1045 }
1046 match options.compression.or_else(|| {
1047 call.paths
1048 .first()
1049 .and_then(|p| CompressionFormat::from_extension(p))
1050 }) {
1051 Some(CompressionFormat::Bzip2 | CompressionFormat::Xz) => Some(Source::Placeholder {
1052 what: format!(
1053 "{names}: Polars cannot read bzip2 or xz; decompress it and read it with pl.scan_csv."
1054 ),
1055 }),
1056 _ => None,
1057 }
1058}
1059
1060pub(crate) fn arrow_arguments(call: &mut Call<'_>) -> Option<Source> {
1063 let options = call.record.options;
1064 let inputs: Option<Vec<(String, bool)>> = match &options.arrow_parts {
1065 Some(parts) => Some(
1066 parts
1067 .iter()
1068 .map(|part| match part {
1069 crate::formats::ipc_stream::Part::InPlace(p) => (p, false),
1070 crate::formats::ipc_stream::Part::Converted { source, .. } => (source, true),
1071 })
1072 .map(|(p, stream)| (without_secrets(&p.to_string_lossy()).0, stream))
1073 .collect(),
1074 ),
1075 None => hugging_face_files(call.paths, options.table.as_deref())
1076 .map(|files| files.into_iter().map(|f| (f, false)).collect()),
1077 };
1078 match inputs {
1079 Some(inputs) => {
1080 let extra = call.storage_for(inputs.iter().map(|(name, _)| name));
1081 let read = arrow_read(&inputs, extra.as_deref());
1082 let mut after = std::mem::take(&mut call.after);
1083 after.extend(call.skip_tail.take());
1084 Some(Source::Read {
1085 call: read,
1086 after,
1087 notes: std::mem::take(&mut call.notes),
1088 imports: Vec::new(),
1089 })
1090 }
1091 None => {
1092 if let Some(s) = call.storage_for(call.names) {
1093 call.args.push(s);
1094 }
1095 None
1096 }
1097 }
1098}
1099
1100pub(crate) fn excel_arguments(call: &mut Call<'_>) -> Option<Source> {
1102 if let Some(sheet) = &call.record.options.table {
1103 match sheet.parse::<usize>() {
1104 Ok(i) => call.args.push(format!("sheet_id={}", i + 1)),
1106 Err(_) => call.args.push(format!("sheet_name={}", py_str(sheet))),
1107 }
1108 }
1109 call.notes.push(
1110 "datui types a worksheet's columns itself; Polars may read some differently.".to_string(),
1111 );
1112 None
1113}
1114
1115fn sql_ident(name: &str) -> String {
1117 format!("\"{}\"", name.replace('"', "\"\""))
1118}
1119
1120fn whole(call: &mut Call<'_>, read: String, imports: Vec<&'static str>) -> Source {
1123 call.notes.extend(read_whole_note(call.record, &read));
1124 let mut after = std::mem::take(&mut call.after);
1125 after.extend(call.skip_tail.take());
1126 Source::Read {
1127 call: read,
1128 after,
1129 notes: std::mem::take(&mut call.notes),
1130 imports,
1131 }
1132}
1133
1134fn read_whole_note(record: &OpenRecord, call: &str) -> Option<String> {
1137 let name = call.split('(').next().unwrap_or(call);
1138 (record.read_mode == Some(crate::ReadMode::Lazy)).then(|| {
1139 format!(
1140 "Read: {} in datui; {name} reads the file whole into memory.",
1141 crate::ReadMode::Lazy.label()
1142 )
1143 })
1144}
1145
1146pub(crate) fn sqlite_arguments(call: &mut Call<'_>) -> Option<Source> {
1148 let [file] = call.names else {
1149 return None;
1150 };
1151 let Some(table) = call.record.options.table.as_deref() else {
1152 return Some(Source::Placeholder {
1153 what: format!("{file}: datui could not tell which table it read; load it here."),
1154 });
1155 };
1156 if call.paths.iter().any(|p| is_url(p)) {
1157 return Some(Source::Placeholder {
1158 what: format!(
1159 "{file} --table {table}: sqlite3 opens a local file; download it and read it \
1160 with pl.read_database."
1161 ),
1162 });
1163 }
1164 call.notes.push(
1165 "datui types a table's columns from their declared types; Polars infers them from \
1166 the values."
1167 .to_string(),
1168 );
1169 let query = format!("SELECT * FROM {}", sql_ident(table));
1170 let read = format!(
1171 "pl.read_database({}, sqlite3.connect({})).lazy()",
1172 py_str(&query),
1173 py_str(file)
1174 );
1175 Some(whole(call, read, vec!["import sqlite3"]))
1176}
1177
1178pub(crate) fn numpy_arguments(call: &mut Call<'_>) -> Option<Source> {
1180 let [file] = call.names else {
1181 return None;
1182 };
1183 let table = call.record.options.table.as_deref();
1184 let archive = table.is_some() || file.to_ascii_lowercase().ends_with(".npz");
1185 let array = match table {
1186 Some(name) => format!("np.load({})[{}]", py_str(file), py_str(name)),
1187 None if archive => format!("next(iter(np.load({}).values()))", py_str(file)),
1189 None => format!("np.load({})", py_str(file)),
1190 };
1191 let names: Vec<String> = call
1192 .record
1193 .schema
1194 .iter_names()
1195 .map(|n| n.to_string())
1196 .collect();
1197 let schema = if names.iter().any(|n| n.contains('.')) {
1199 call.notes.push(
1200 "datui names a nested field's columns outer.inner; Polars keeps the field as a struct."
1201 .to_string(),
1202 );
1203 String::new()
1204 } else {
1205 format!(", schema={}", py_names(&names))
1206 };
1207 let read = format!("pl.from_numpy({array}{schema}, orient=\"row\").lazy()");
1208 Some(whole(call, read, vec!["import numpy as np"]))
1209}
1210
1211fn named_with_table(names: &[String], record: &OpenRecord) -> String {
1214 let names = names.join(", ");
1215 match record.options.table.as_deref() {
1216 Some(table) => format!("{names} --table {table}"),
1217 None => names,
1218 }
1219}
1220
1221pub(crate) fn lines_arguments(call: &mut Call<'_>) -> Option<Source> {
1224 let names = call.names.join(", ");
1225 let path = match call.paths {
1226 [one]
1227 if !is_url(one)
1228 && !one.is_dir()
1229 && CompressionFormat::from_extension(one).is_none() =>
1230 {
1231 one
1232 }
1233 _ => {
1234 return Some(Source::Placeholder {
1235 what: format!("{names}: datui read it as lines; load it here."),
1236 });
1237 }
1238 };
1239 let read = format!(
1240 "pl.LazyFrame({{\"line\": open({}, encoding=\"utf-8\", errors=\"replace\", newline=\"\").read().removesuffix(\"\\n\").split(\"\\n\")}})",
1241 py_str(&path.to_string_lossy())
1242 );
1243 let mut after = vec![".with_columns(pl.col(\"line\").str.strip_suffix(\"\\r\"))".to_string()];
1244 after.append(&mut call.after);
1245 Some(Source::Read {
1246 call: read,
1247 after,
1248 notes: std::mem::take(&mut call.notes),
1249 imports: Vec::new(),
1250 })
1251}
1252
1253pub fn source(record: &OpenRecord) -> Source {
1255 let Some(paths) = record.paths else {
1256 return Source::Placeholder {
1257 what: "The data datui was handed: load it here as a DataFrame or LazyFrame."
1258 .to_string(),
1259 };
1260 };
1261 let teed;
1263 let paths = match (&record.options.tee, paths) {
1264 (Some(tee), [one]) if crate::loading::stdin::is_stdin(one) => {
1265 teed = [tee.clone()];
1266 &teed[..]
1267 }
1268 _ => paths,
1269 };
1270 if paths.iter().any(|p| crate::loading::stdin::is_stdin(p)) {
1271 return Source::Placeholder {
1272 what: "The data datui read from standard input: load it here.".to_string(),
1273 };
1274 }
1275 let spec = record.spec.clone().or_else(|| {
1276 let options = record.options;
1277 options
1278 .spec_name
1279 .clone()
1280 .or_else(|| options.spec_file.as_ref().map(|f| f.display().to_string()))
1281 });
1282 if let Some(spec) = spec {
1283 return Source::Placeholder {
1284 what: format!(
1285 "datui read this through the format spec {spec}, which it cannot write as \
1286 Python: load it here."
1287 ),
1288 };
1289 }
1290 let targets: Option<Vec<Target>> = paths.iter().map(|p| reader_target(p, record)).collect();
1291 let Some(targets) = targets.filter(|t| !t.is_empty()) else {
1292 let names: Vec<String> = paths
1293 .iter()
1294 .map(|p| without_secrets(&p.to_string_lossy()).0)
1295 .collect();
1296 return Source::Placeholder {
1297 what: format!(
1298 "{}: datui could not name a Polars reader for this data; load it here.",
1299 named_with_table(&names, record)
1300 ),
1301 };
1302 };
1303 let format = targets[0].format;
1304 if targets.iter().any(|t| t.format != format) {
1305 return Source::Placeholder {
1306 what: "The files are of more than one format: load them here.".to_string(),
1307 };
1308 }
1309 let below = targets.iter().any(|t| t.below);
1310 let python = python_of(format);
1311 let literal = targets.iter().any(|t| t.literal);
1314 let no_glob = literal
1315 && python.is_some_and(|python| python.glob_flag)
1316 && !targets.iter().any(|t| t.pattern);
1317 let names: Vec<String> = targets
1318 .into_iter()
1319 .map(|t| {
1320 if t.literal && !no_glob {
1321 crate::cloud::source::escape_glob(&t.text)
1322 } else {
1323 t.text
1324 }
1325 })
1326 .collect();
1327 let Some(python) = python else {
1328 return Source::Placeholder {
1329 what: format!(
1330 "{}: Polars has no reader for {} files; load it here.",
1331 named_with_table(&names, record),
1332 format.title()
1333 ),
1334 };
1335 };
1336 let target = match names.as_slice() {
1337 [one] => py_str(one),
1338 many => py_names(many),
1339 };
1340 let options = record.options;
1341 let mut args: Vec<String> = vec![target];
1342 if no_glob {
1343 args.push("glob=False".to_string());
1344 }
1345 let mut notes = Vec::new();
1346 let endpoint = record.s3_endpoint.as_deref().map(without_secrets);
1347 if paths
1348 .iter()
1349 .any(|p| is_url(p) && without_secrets(&p.to_string_lossy()).1)
1350 || endpoint.as_ref().is_some_and(|(_, cut)| *cut)
1351 {
1352 notes.push(
1353 "datui left a user, password or query string out of the URL, as it may be a \
1354 credential: add it back if the server needs it."
1355 .to_string(),
1356 );
1357 }
1358 let in_store = names.iter().find(|n| store_scheme(n).is_some());
1359 let storage = in_store
1360 .and_then(|name| storage_options(name, record, endpoint.as_ref().map(|(e, _)| e.as_str())));
1361 if let Some(name) = in_store
1363 && python.eager
1364 {
1365 notes.push(format!(
1366 "{} reads no object store: download {name} and read it from disk.",
1367 python.call
1368 ));
1369 }
1370 let mut call = Call {
1371 record,
1372 paths,
1373 format,
1374 names: &names,
1375 below,
1376 storage,
1377 args,
1378 after: options.read_python.clone(),
1380 skip_tail: None,
1381 notes,
1382 };
1383 if let Some(arguments) = python.arguments
1384 && let Some(source) = arguments(&mut call)
1385 {
1386 return source;
1387 }
1388 let Call {
1389 args,
1390 mut after,
1391 skip_tail,
1392 mut notes,
1393 ..
1394 } = call;
1395 after.extend(skip_tail);
1396 if !record.read_as_text.is_empty() {
1397 notes.push(format!(
1398 "datui read these columns as text from every file: {}.",
1399 record.read_as_text.join(", ")
1400 ));
1401 }
1402 let mut call = format!("{}({})", python.call, args.join(", "));
1403 if python.eager {
1404 notes.extend(read_whole_note(record, &call));
1405 call.push_str(".lazy()");
1406 }
1407 Source::Read {
1408 call,
1409 after,
1410 notes,
1411 imports: Vec::new(),
1412 }
1413}
1414
1415fn csv_null_values(options: &OpenOptions, schema: &Schema) -> Option<String> {
1419 let specs = options.null_values.as_ref().filter(|s| !s.is_empty())?;
1420 let mut global = Vec::new();
1421 let mut per_column: Vec<(String, String)> = Vec::new();
1422 for spec in specs {
1423 match spec.find('=') {
1424 Some(i) => per_column.push((spec[..i].to_string(), spec[i + 1..].to_string())),
1425 None => global.push(spec.clone()),
1426 }
1427 }
1428 let dict = |pairs: Vec<(String, String)>| {
1429 let items: Vec<String> = pairs
1430 .iter()
1431 .map(|(c, v)| format!("{}: {}", py_str(c), py_str(v)))
1432 .collect();
1433 format!("{{{}}}", items.join(", "))
1434 };
1435 Some(match (global.as_slice(), per_column.is_empty()) {
1436 ([one], true) => py_str(one),
1437 (_, true) => py_names(&global),
1438 ([], false) => dict(per_column),
1439 (_, false) => dict(
1440 schema
1441 .iter_names()
1442 .map(|name| {
1443 let value = per_column
1444 .iter()
1445 .rev()
1446 .find(|(c, _)| c == name.as_str())
1447 .map(|(_, v)| v.clone())
1448 .unwrap_or_else(|| global[0].clone());
1449 (name.to_string(), value)
1450 })
1451 .collect(),
1452 ),
1453 })
1454}
1455
1456pub fn py_value(value: &AnyValue) -> Option<String> {
1459 Some(match value {
1460 AnyValue::Null => "None".to_string(),
1461 AnyValue::Boolean(b) => py_bool(*b).to_string(),
1462 AnyValue::String(s) => py_str(s),
1463 AnyValue::StringOwned(s) => py_str(s),
1464 AnyValue::Int8(v) => v.to_string(),
1465 AnyValue::Int16(v) => v.to_string(),
1466 AnyValue::Int32(v) => v.to_string(),
1467 AnyValue::Int64(v) => v.to_string(),
1468 AnyValue::UInt8(v) => v.to_string(),
1469 AnyValue::UInt16(v) => v.to_string(),
1470 AnyValue::UInt32(v) => v.to_string(),
1471 AnyValue::UInt64(v) => v.to_string(),
1472 AnyValue::Float32(v) => py_float(f64::from(*v)),
1473 AnyValue::Float64(v) => py_float(*v),
1474 AnyValue::Date(days) => {
1475 let date = chrono::NaiveDate::from_ymd_opt(1970, 1, 1)?
1476 .checked_add_signed(chrono::Duration::days(i64::from(*days)))?;
1477 use chrono::Datelike;
1478 format!("pl.date({}, {}, {})", date.year(), date.month(), date.day())
1479 }
1480 _ => return None,
1481 })
1482}
1483
1484#[derive(Debug, Clone, PartialEq)]
1486pub struct Script {
1487 pub source: Source,
1488 pub steps: Vec<Step>,
1489}
1490
1491impl Script {
1492 pub fn render(&self) -> String {
1493 let mut out = String::from("import polars as pl\n");
1494 if let Source::Read { imports, .. } = &self.source {
1495 for import in imports {
1496 out.push_str(import);
1497 out.push('\n');
1498 }
1499 }
1500 out.push('\n');
1501 let (head, mut lines) = match &self.source {
1502 Source::Read {
1503 call, after, notes, ..
1504 } => {
1505 for note in notes {
1506 out.push_str(&py_comment(note));
1507 out.push('\n');
1508 }
1509 (call.clone(), after.clone())
1510 }
1511 Source::Placeholder { what } => {
1512 out.push_str(&py_comment(what));
1513 out.push_str("\ndf = ...\n\n");
1514 ("df.lazy()".to_string(), Vec::new())
1515 }
1516 };
1517 let mut stopped = false;
1520 for step in &self.steps {
1521 let calls = step.python();
1522 if stopped {
1523 for call in &calls {
1525 lines.extend(call.lines().map(|c| {
1526 if c.starts_with('#') {
1527 c.to_string()
1528 } else {
1529 format!("# {c}")
1530 }
1531 }));
1532 }
1533 } else {
1534 stopped = matches!(step, Step::Unreproducible(_));
1535 lines.extend(calls);
1536 }
1537 }
1538 if lines.is_empty() {
1539 out.push_str(&format!("df = {head}\n"));
1540 } else {
1541 out.push_str("df = (\n");
1542 out.push_str(&format!(" {head}\n"));
1543 for line in lines {
1544 out.push_str(&format!(" {line}\n"));
1545 }
1546 out.push_str(")\n");
1547 }
1548 out
1549 }
1550}
1551
1552#[cfg(test)]
1553mod tests;