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_time_zones(input);
426 if let Err(e) = nodes.resolve_types(input) {
427 return vec![py_comment(&format!("the query did not run: {e}"))];
428 }
429 nodes.python_steps(keys)
430 }
431 Err(e) => vec![py_comment(&format!("the query did not parse: {e}"))],
432 },
433 Step::QueryRows { query, input } => match crate::query::parse_nodes(query) {
434 Ok(mut nodes) => {
435 nodes.resolve_time_zones(input);
436 if let Err(e) = nodes.resolve_types(input) {
437 return vec![py_comment(&format!("the query did not run: {e}"))];
438 }
439 nodes.python_filters()
440 }
441 Err(e) => vec![py_comment(&format!("the query did not parse: {e}"))],
442 },
443 Step::Sql { sql, ordered_by } => {
444 let sql = sql.trim();
445 let verbatim = sql.contains('\n')
448 && !sql.contains("\"\"\"")
449 && !sql.ends_with('"')
450 && sql
451 .chars()
452 .all(|c| c == '\n' || c == '\t' || (c != '\\' && !c.is_control()));
453 let mut lines = if verbatim {
454 vec![
455 ".sql(".to_string(),
456 format!(" \"\"\"{sql}\"\"\","),
457 " table_name=\"df\",".to_string(),
458 ")".to_string(),
459 ]
460 } else {
461 vec![format!(".sql({}, table_name=\"df\")", py_str(sql))]
462 };
463 if !ordered_by.is_empty() {
464 lines.push(sort_call(ordered_by, &vec![false; ordered_by.len()]));
465 }
466 lines
467 }
468 Step::Search { patterns, columns } => {
469 let terms: Vec<String> = patterns
470 .iter()
471 .map(|p| {
472 let any: Vec<String> = columns
473 .iter()
474 .map(|c| {
475 format!(
476 "pl.col({}).str.contains({}, strict=False)",
477 py_str(c),
478 py_str(p)
479 )
480 })
481 .collect();
482 if any.len() == 1 || patterns.len() == 1 {
483 any.join(" | ")
484 } else {
485 format!("({})", any.join(" | "))
486 }
487 })
488 .collect();
489 vec![format!(".filter({})", terms.join(" & "))]
490 }
491 Step::Filter(filters) => vec![format!(".filter({})", filters_python(filters))],
492 Step::Sort {
493 columns,
494 descending,
495 } => vec![sort_call(columns, descending)],
496 Step::Reverse => vec![".reverse()".to_string()],
497 Step::Select(columns) => vec![format!(".select({})", py_names(columns))],
498 Step::Drop(columns) => vec![format!(".drop({})", py_names(columns))],
499 Step::Pivot {
500 index,
501 on,
502 values,
503 aggregation,
504 } => {
505 let agg = match aggregation {
508 PivotAggregation::Last => "last()",
509 PivotAggregation::First => "first()",
510 PivotAggregation::Min => "min()",
511 PivotAggregation::Max => "max()",
512 PivotAggregation::Avg => "mean()",
513 PivotAggregation::Med => "median()",
514 PivotAggregation::Std => "std()",
515 PivotAggregation::Count => "len()",
516 };
517 let cell = match aggregation {
518 PivotAggregation::Count => "sum",
519 _ => "first",
520 };
521 let keys: Vec<String> = index.iter().chain([on]).cloned().collect();
522 vec![
523 format!(".group_by({}, maintain_order=True)", py_names(&keys)),
524 format!(".agg(pl.col({}).{agg})", py_str(values)),
525 ".collect()".to_string(),
526 ".pipe(".to_string(),
527 " lambda cells: cells.pivot(".to_string(),
528 format!(" on={},", py_str(on)),
529 format!(
530 " on_columns=cells[{}].unique().sort(nulls_last=True),",
531 py_str(on)
532 ),
533 format!(" index={},", py_names(index)),
534 format!(" values={},", py_str(values)),
535 format!(" aggregate_function={},", py_str(cell)),
536 " )".to_string(),
537 ")".to_string(),
538 ".lazy()".to_string(),
539 ]
540 }
541 Step::Melt {
542 index,
543 on,
544 variable_name,
545 value_name,
546 } => vec![format!(
547 ".unpivot(on={}, index={}, variable_name={}, value_name={})",
548 py_names(on),
549 py_names(index),
550 py_str(variable_name),
551 py_str(value_name)
552 )],
553 Step::Matching(keys) => {
554 let terms: Vec<String> = keys
555 .iter()
556 .map(|(key, value)| format!("{key}.eq_missing({value})"))
557 .collect();
558 let terms = if terms.len() == 1 {
559 terms
560 } else {
561 terms.into_iter().map(|t| format!("({t})")).collect()
562 };
563 vec![format!(".filter({})", terms.join(" & "))]
564 }
565 Step::Unreproducible(what) => vec![py_comment(what)],
566 }
567 }
568}
569
570#[derive(Debug, Clone, PartialEq)]
572pub enum Source {
573 Read {
577 call: String,
578 after: Vec<String>,
579 notes: Vec<String>,
580 imports: Vec<&'static str>,
582 },
583 Placeholder { what: String },
586}
587
588pub struct OpenRecord<'a> {
590 pub paths: Option<&'a [PathBuf]>,
592 pub options: &'a OpenOptions,
593 pub format: Option<FileFormat>,
597 pub read_mode: Option<crate::ReadMode>,
599 pub schema: &'a Schema,
601 pub remote_objects: Vec<String>,
603 pub s3_endpoint: Option<String>,
606 pub s3_region: Option<String>,
607 pub unsigned: bool,
609 pub read_as_text: Vec<String>,
611 pub spec: Option<String>,
613}
614
615fn is_url(path: &Path) -> bool {
616 crate::cloud::source::is_remote_url(path)
617}
618
619fn without_secrets(url: &str) -> (String, bool) {
622 let url = &*crate::cloud::source::split_source_id(url).1;
624 let Some(scheme_end) = url.find("://").map(|i| i + 3) else {
625 return (url.to_string(), false);
626 };
627 let (scheme, rest) = url.split_at(scheme_end);
628 let host_end = rest.find(['/', '?', '#']).unwrap_or(rest.len());
629 let (authority, path) = rest.split_at(host_end);
630 let azure = ["abfs://", "abfss://"]
632 .iter()
633 .any(|s| scheme.eq_ignore_ascii_case(s));
634 let host = match authority.rsplit_once('@') {
635 Some((_, host)) if !azure => host,
636 _ => authority,
637 };
638 let http = scheme.eq_ignore_ascii_case("http://") || scheme.eq_ignore_ascii_case("https://");
639 let path = match path.find(['?', '#']) {
640 Some(i) if http => &path[..i],
641 _ => path,
642 };
643 let kept = format!("{scheme}{host}{path}");
644 let cut = kept != url;
645 (kept, cut)
646}
647
648fn file_format(path: &Path, record: &OpenRecord) -> Option<FileFormat> {
651 record.format.or(record.options.format).or_else(|| {
652 FileFormat::from_path(path).or_else(|| {
653 CompressionFormat::from_extension(path)
654 .and_then(|_| path.file_stem())
655 .and_then(|stem| FileFormat::from_path(Path::new(stem)))
656 })
657 })
658}
659
660fn commonest_format<'a>(names: impl Iterator<Item = &'a str>) -> Option<FileFormat> {
662 let mut counts: Vec<(FileFormat, usize)> = Vec::new();
663 for name in names {
664 if let Some(format) = FileFormat::from_path(Path::new(name)) {
665 match counts.iter_mut().find(|(f, _)| *f == format) {
666 Some((_, n)) => *n += 1,
667 None => counts.push((format, 1)),
668 }
669 }
670 }
671 counts.into_iter().max_by_key(|(_, n)| *n).map(|(f, _)| f)
672}
673
674struct Target {
676 text: String,
679 format: FileFormat,
680 below: bool,
682 pattern: bool,
684 literal: bool,
686}
687
688fn reader_target(path: &Path, record: &OpenRecord) -> Option<Target> {
690 let text = path.to_string_lossy().to_string();
691 if is_url(path) {
692 let (text, _) = without_secrets(&text);
693 if let Some(format) = file_format(Path::new(&text), record) {
694 let pattern = crate::cloud::source::has_glob_chars(Path::new(&text));
695 return Some(Target {
696 text,
697 format,
698 below: false,
699 pattern,
700 literal: false,
701 });
702 }
703 let format = record
705 .format
706 .or(record.options.format)
707 .or_else(|| commonest_format(record.remote_objects.iter().map(String::as_str)))?;
708 let base = text.trim_end_matches('/');
709 let ext = format_extension(format)?;
710 return Some(Target {
711 text: format!("{base}/**/*.{ext}"),
712 format,
713 below: true,
714 pattern: true,
715 literal: false,
716 });
717 }
718 if path.is_dir() {
719 let entries: Vec<std::fs::DirEntry> = std::fs::read_dir(path).ok()?.flatten().collect();
720 let mut names: Vec<String> = entries
721 .iter()
722 .filter(|e| e.path().is_file())
723 .map(|e| e.file_name().to_string_lossy().to_string())
724 .collect();
725 if names
727 .iter()
728 .any(|n| FileFormat::from_path(Path::new(n)) == Some(FileFormat::Arrow))
729 {
730 names.retain(|n| !crate::home::discover::is_hugging_face_metadata(n));
731 }
732 let has_dirs = entries.iter().any(|e| e.path().is_dir());
733 let format = record
734 .format
735 .or(record.options.format)
736 .or_else(|| commonest_format(names.iter().map(String::as_str)));
737 let base = crate::cloud::source::escape_glob(text.trim_end_matches(['/', '\\']));
739 let (text, format, below) = match format {
742 Some(FileFormat::Parquet) | None if has_dirs => {
743 (format!("{base}/**/*.parquet"), FileFormat::Parquet, true)
744 }
745 Some(format) => (
746 format!("{base}/*.{}", format_extension(format)?),
747 format,
748 false,
749 ),
750 None => return None,
751 };
752 return Some(Target {
753 text,
754 format,
755 below,
756 pattern: true,
757 literal: false,
758 });
759 }
760 let format = file_format(path, record)?;
761 let pattern = crate::cloud::source::expands_as_glob(path);
762 Some(Target {
763 literal: !pattern && scans_by_pattern(format) && crate::cloud::source::has_glob_chars(path),
766 text,
767 format,
768 below: false,
769 pattern,
770 })
771}
772
773fn scans_by_pattern(format: FileFormat) -> bool {
775 python_of(format).is_some_and(|python| !python.eager)
776}
777
778fn hugging_face_files(paths: &[PathBuf], table: Option<&str>) -> Option<Vec<String>> {
781 let [path] = paths else {
782 return None;
783 };
784 if is_url(path) || !path.is_dir() {
785 return None;
786 }
787 let dict_split = crate::formats::hf_splits::dataset_dict(path).and_then(|splits| {
789 let listed: Vec<&str> = splits.iter().map(String::as_str).collect();
790 crate::formats::hf_splits::pick(&listed, table).ok()?.split
791 });
792 let dict = dict_split.is_some();
793 let table = if dict { None } else { table };
794 let path = &dict_split.map_or_else(|| path.clone(), |split| path.join(split));
795 let mut inside: Vec<PathBuf> = std::fs::read_dir(path)
796 .ok()?
797 .flatten()
798 .map(|e| e.path())
799 .filter(|p| p.is_file() && FileFormat::from_path(p) == Some(FileFormat::Arrow))
800 .collect();
801 inside.sort();
802 let cache = ["dataset_info.json", "state.json"]
803 .iter()
804 .any(|name| path.join(name).is_file());
805 if !cache && !dict {
806 return None;
807 }
808 if cache {
809 let names: Vec<&str> = inside
810 .iter()
811 .map(|f| f.file_name().and_then(|n| n.to_str()).unwrap_or_default())
812 .collect();
813 let (chosen, _) = crate::formats::hf_splits::choose(&names, table).ok()?;
814 inside = chosen.into_iter().map(|i| inside[i].clone()).collect();
815 }
816 Some(
817 inside
818 .iter()
819 .map(|p| p.to_string_lossy().to_string())
820 .collect(),
821 )
822}
823
824fn arrow_read(inputs: &[(String, bool)], extra: Option<&str>) -> String {
828 let extra = extra.map(|e| format!(", {e}")).unwrap_or_default();
829 let names: Vec<String> = inputs.iter().map(|(name, _)| name.clone()).collect();
830 if inputs.iter().all(|(_, stream)| *stream) {
831 return match names.as_slice() {
832 [one] => format!("pl.read_ipc_stream({}{extra}).lazy()", py_str(one)),
833 many => format!(
834 "pl.concat([pl.read_ipc_stream(f{extra}) for f in {}]).lazy()",
835 py_names(many)
836 ),
837 };
838 }
839 if inputs.iter().all(|(_, stream)| !*stream) {
840 return match names.as_slice() {
841 [one] => format!("pl.scan_ipc({}{extra})", py_str(one)),
842 many => format!("pl.scan_ipc({}{extra})", py_names(many)),
843 };
844 }
845 let reads: Vec<String> = inputs
846 .iter()
847 .map(|(name, stream)| match stream {
848 true => format!("pl.read_ipc_stream({}{extra}).lazy()", py_str(name)),
849 false => format!("pl.scan_ipc({}{extra})", py_str(name)),
850 })
851 .collect();
852 format!(
853 "pl.concat([{}], how=\"diagonal_relaxed\")",
854 reads.join(", ")
855 )
856}
857
858fn format_extension(format: FileFormat) -> Option<&'static str> {
861 python_of(format)?;
862 let d = format.descriptor();
863 (d.many_files || format.separator().is_some())
864 .then(|| d.extensions.first().copied())
865 .flatten()
866}
867
868pub(crate) struct Python {
871 pub call: &'static str,
873 pub eager: bool,
875 pub glob_flag: bool,
877 pub arguments: Option<fn(&mut Call<'_>) -> Option<Source>>,
880}
881
882pub(crate) struct Call<'a> {
884 pub record: &'a OpenRecord<'a>,
885 pub paths: &'a [PathBuf],
886 pub format: FileFormat,
887 pub names: &'a [String],
889 pub below: bool,
891 pub storage: Option<String>,
893 pub args: Vec<String>,
894 pub after: Vec<String>,
895 pub skip_tail: Option<String>,
897 pub notes: Vec<String>,
898}
899
900impl Call<'_> {
901 fn storage_for(&self, names: impl IntoIterator<Item = impl AsRef<str>>) -> Option<String> {
903 names
904 .into_iter()
905 .any(|n| store_scheme(n.as_ref()).is_some())
906 .then(|| self.storage.clone())
907 .flatten()
908 }
909}
910
911fn store_scheme(name: &str) -> Option<&'static str> {
914 let (scheme, _) = name.split_once("://")?;
915 match scheme.to_ascii_lowercase().as_str() {
916 "s3" | "s3a" => Some("s3"),
917 "gs" | "gcs" => Some("gs"),
918 "az" | "adl" | "azure" | "abfs" | "abfss" => Some("azure"),
919 _ => None,
920 }
921}
922
923fn storage_options(name: &str, record: &OpenRecord, endpoint: Option<&str>) -> Option<String> {
927 let mut pairs: Vec<(&str, String)> = Vec::new();
928 match store_scheme(name)? {
929 "s3" => {
930 pairs.extend(endpoint.map(|e| ("aws_endpoint_url", e.to_string())));
931 pairs.extend(record.s3_region.clone().map(|r| ("aws_region", r)));
932 }
933 "azure" => {
934 pairs.extend(
935 crate::cloud::source::azure_parts(name)
936 .map(|(account, ..)| ("account_name", account)),
937 );
938 }
939 _ => {}
940 }
941 if record.unsigned {
942 pairs.push(("skip_signature", "true".to_string()));
943 }
944 let pairs: Vec<String> = pairs
945 .iter()
946 .map(|(k, v)| format!("{}: {}", py_str(k), py_str(v)))
947 .collect();
948 (!pairs.is_empty()).then(|| format!("storage_options={{{}}}", pairs.join(", ")))
949}
950
951pub(crate) fn ndjson_arguments(call: &mut Call<'_>) -> Option<Source> {
953 if let Some(s) = call.storage_for(call.names) {
954 call.args.push(s);
955 }
956 None
957}
958
959fn python_of(format: FileFormat) -> Option<&'static Python> {
960 crate::formats::readers::of(format).python.as_ref()
961}
962
963pub(crate) fn parquet_arguments(call: &mut Call<'_>) -> Option<Source> {
965 if call.record.options.hive || call.below {
966 call.args.push("hive_partitioning=True".to_string());
967 }
968 if let Some(s) = call.storage_for(call.names) {
969 call.args.push(s);
970 }
971 None
972}
973
974pub(crate) fn csv_arguments(call: &mut Call<'_>) -> Option<Source> {
976 let options = call.record.options;
977 let names = call.names.join(", ");
978 let comment = options.comment_char.as_deref().filter(|c| !c.is_empty());
980 let dialect: Vec<&str> = [
981 (comment.is_some_and(|c| c.len() > 5), "--comment"),
982 (options.header_rows().is_some(), "--header-rows"),
983 (options.skip_initial_space, "--skip-initial-space"),
984 ]
985 .into_iter()
986 .filter_map(|(set, flag)| set.then_some(flag))
987 .collect();
988 if !dialect.is_empty() {
989 return Some(Source::Placeholder {
990 what: format!(
991 "{names}: datui read it with {}, which it cannot write as Python: load it here.",
992 dialect.join(", ")
993 ),
994 });
995 }
996 let separator = options
997 .delimiter
998 .or_else(|| call.format.separator())
999 .unwrap_or(b',');
1000 let args = &mut call.args;
1001 if separator != b',' {
1002 args.push(format!(
1003 "separator={}",
1004 py_str(&(separator as char).to_string())
1005 ));
1006 }
1007 if let Some(prefix) = comment {
1008 args.push(format!("comment_prefix={}", py_str(prefix)));
1009 }
1010 if options.has_header == Some(false) {
1011 args.push("has_header=False".to_string());
1012 }
1013 if let Some(n) = options.skip_lines {
1014 args.push(format!("skip_lines={n}"));
1015 }
1016 if let Some(n) = options.skip_rows {
1017 args.push(format!("skip_rows={n}"));
1018 }
1019 if let Some(n) = options.infer_schema_length {
1020 args.push(format!("infer_schema_length={n}"));
1021 }
1022 if options.ignore_errors {
1023 args.push("ignore_errors=True".to_string());
1024 }
1025 let bucket_prefix = call.below && call.paths.iter().any(|p| is_url(p));
1027 if options.csv_try_parse_dates() && !bucket_prefix {
1028 call.args.push("try_parse_dates=True".to_string());
1029 }
1030 if let Some(nulls) = csv_null_values(options, call.record.schema) {
1031 call.args.push(format!("null_values={nulls}"));
1032 }
1033 if !options.typing.text.is_empty() {
1035 let text: Vec<String> = options
1036 .typing
1037 .text
1038 .iter()
1039 .map(|name| format!("{}: pl.String", py_str(name)))
1040 .collect();
1041 call.args
1042 .push(format!("schema_overrides={{{}}}", text.join(", ")));
1043 }
1044 if let Some(s) = call.storage_for(call.names) {
1045 call.args.push(s);
1046 }
1047 if let Some(n) = options.skip_tail_rows.filter(|n| *n > 0) {
1048 call.skip_tail = Some(format!(".filter(pl.int_range(pl.len()) < pl.len() - {n})"));
1049 }
1050 match options.compression.or_else(|| {
1051 call.paths
1052 .first()
1053 .and_then(|p| CompressionFormat::from_extension(p))
1054 }) {
1055 Some(CompressionFormat::Bzip2 | CompressionFormat::Xz) => Some(Source::Placeholder {
1056 what: format!(
1057 "{names}: Polars cannot read bzip2 or xz; decompress it and read it with pl.scan_csv."
1058 ),
1059 }),
1060 _ => None,
1061 }
1062}
1063
1064pub(crate) fn arrow_arguments(call: &mut Call<'_>) -> Option<Source> {
1067 let options = call.record.options;
1068 let inputs: Option<Vec<(String, bool)>> = match &options.arrow_parts {
1069 Some(parts) => Some(
1070 parts
1071 .iter()
1072 .map(|part| match part {
1073 crate::formats::ipc_stream::Part::InPlace(p) => (p, false),
1074 crate::formats::ipc_stream::Part::Converted { source, .. } => (source, true),
1075 })
1076 .map(|(p, stream)| (without_secrets(&p.to_string_lossy()).0, stream))
1077 .collect(),
1078 ),
1079 None => hugging_face_files(call.paths, options.table.as_deref())
1080 .map(|files| files.into_iter().map(|f| (f, false)).collect()),
1081 };
1082 match inputs {
1083 Some(inputs) => {
1084 let extra = call.storage_for(inputs.iter().map(|(name, _)| name));
1085 let read = arrow_read(&inputs, extra.as_deref());
1086 let mut after = std::mem::take(&mut call.after);
1087 after.extend(call.skip_tail.take());
1088 Some(Source::Read {
1089 call: read,
1090 after,
1091 notes: std::mem::take(&mut call.notes),
1092 imports: Vec::new(),
1093 })
1094 }
1095 None => {
1096 if let Some(s) = call.storage_for(call.names) {
1097 call.args.push(s);
1098 }
1099 None
1100 }
1101 }
1102}
1103
1104pub(crate) fn excel_arguments(call: &mut Call<'_>) -> Option<Source> {
1106 if let Some(sheet) = &call.record.options.table {
1107 match sheet.parse::<usize>() {
1108 Ok(i) => call.args.push(format!("sheet_id={}", i + 1)),
1110 Err(_) => call.args.push(format!("sheet_name={}", py_str(sheet))),
1111 }
1112 }
1113 call.notes.push(
1114 "datui types a worksheet's columns itself; Polars may read some differently.".to_string(),
1115 );
1116 None
1117}
1118
1119fn sql_ident(name: &str) -> String {
1121 format!("\"{}\"", name.replace('"', "\"\""))
1122}
1123
1124fn whole(call: &mut Call<'_>, read: String, imports: Vec<&'static str>) -> Source {
1127 call.notes.extend(read_whole_note(call.record, &read));
1128 let mut after = std::mem::take(&mut call.after);
1129 after.extend(call.skip_tail.take());
1130 Source::Read {
1131 call: read,
1132 after,
1133 notes: std::mem::take(&mut call.notes),
1134 imports,
1135 }
1136}
1137
1138fn read_whole_note(record: &OpenRecord, call: &str) -> Option<String> {
1141 let name = call.split('(').next().unwrap_or(call);
1142 (record.read_mode == Some(crate::ReadMode::Lazy)).then(|| {
1143 format!(
1144 "Read: {} in datui; {name} reads the file whole into memory.",
1145 crate::ReadMode::Lazy.label()
1146 )
1147 })
1148}
1149
1150pub(crate) fn sqlite_arguments(call: &mut Call<'_>) -> Option<Source> {
1152 let [file] = call.names else {
1153 return None;
1154 };
1155 let Some(table) = call.record.options.table.as_deref() else {
1156 return Some(Source::Placeholder {
1157 what: format!("{file}: datui could not tell which table it read; load it here."),
1158 });
1159 };
1160 if call.paths.iter().any(|p| is_url(p)) {
1161 return Some(Source::Placeholder {
1162 what: format!(
1163 "{file} --table {table}: sqlite3 opens a local file; download it and read it \
1164 with pl.read_database."
1165 ),
1166 });
1167 }
1168 call.notes.push(
1169 "datui types a table's columns from their declared types; Polars infers them from \
1170 the values."
1171 .to_string(),
1172 );
1173 let query = format!("SELECT * FROM {}", sql_ident(table));
1174 let read = format!(
1175 "pl.read_database({}, sqlite3.connect({})).lazy()",
1176 py_str(&query),
1177 py_str(file)
1178 );
1179 Some(whole(call, read, vec!["import sqlite3"]))
1180}
1181
1182pub(crate) fn numpy_arguments(call: &mut Call<'_>) -> Option<Source> {
1184 let [file] = call.names else {
1185 return None;
1186 };
1187 let table = call.record.options.table.as_deref();
1188 let archive = table.is_some() || file.to_ascii_lowercase().ends_with(".npz");
1189 let array = match table {
1190 Some(name) => format!("np.load({})[{}]", py_str(file), py_str(name)),
1191 None if archive => format!("next(iter(np.load({}).values()))", py_str(file)),
1193 None => format!("np.load({})", py_str(file)),
1194 };
1195 let names: Vec<String> = call
1196 .record
1197 .schema
1198 .iter_names()
1199 .map(|n| n.to_string())
1200 .collect();
1201 let schema = if names.iter().any(|n| n.contains('.')) {
1203 call.notes.push(
1204 "datui names a nested field's columns outer.inner; Polars keeps the field as a struct."
1205 .to_string(),
1206 );
1207 String::new()
1208 } else {
1209 format!(", schema={}", py_names(&names))
1210 };
1211 let read = format!("pl.from_numpy({array}{schema}, orient=\"row\").lazy()");
1212 Some(whole(call, read, vec!["import numpy as np"]))
1213}
1214
1215fn named_with_table(names: &[String], record: &OpenRecord) -> String {
1218 let names = names.join(", ");
1219 match record.options.table.as_deref() {
1220 Some(table) => format!("{names} --table {table}"),
1221 None => names,
1222 }
1223}
1224
1225pub(crate) fn lines_arguments(call: &mut Call<'_>) -> Option<Source> {
1228 let names = call.names.join(", ");
1229 let path = match call.paths {
1230 [one]
1231 if !is_url(one)
1232 && !one.is_dir()
1233 && CompressionFormat::from_extension(one).is_none() =>
1234 {
1235 one
1236 }
1237 _ => {
1238 return Some(Source::Placeholder {
1239 what: format!("{names}: datui read it as lines; load it here."),
1240 });
1241 }
1242 };
1243 let read = format!(
1244 "pl.LazyFrame({{\"line\": open({}, encoding=\"utf-8\", errors=\"replace\", newline=\"\").read().removesuffix(\"\\n\").split(\"\\n\")}})",
1245 py_str(&path.to_string_lossy())
1246 );
1247 let mut after = vec![".with_columns(pl.col(\"line\").str.strip_suffix(\"\\r\"))".to_string()];
1248 after.append(&mut call.after);
1249 Some(Source::Read {
1250 call: read,
1251 after,
1252 notes: std::mem::take(&mut call.notes),
1253 imports: Vec::new(),
1254 })
1255}
1256
1257pub fn source(record: &OpenRecord) -> Source {
1259 let Some(paths) = record.paths else {
1260 return Source::Placeholder {
1261 what: "The data datui was handed: load it here as a DataFrame or LazyFrame."
1262 .to_string(),
1263 };
1264 };
1265 let teed;
1267 let paths = match (&record.options.tee, paths) {
1268 (Some(tee), [one]) if crate::loading::stdin::is_stdin(one) => {
1269 teed = [tee.clone()];
1270 &teed[..]
1271 }
1272 _ => paths,
1273 };
1274 if paths.iter().any(|p| crate::loading::stdin::is_stdin(p)) {
1275 return Source::Placeholder {
1276 what: "The data datui read from standard input: load it here.".to_string(),
1277 };
1278 }
1279 let spec = record.spec.clone().or_else(|| {
1280 let options = record.options;
1281 options
1282 .spec_name
1283 .clone()
1284 .or_else(|| options.spec_file.as_ref().map(|f| f.display().to_string()))
1285 });
1286 if let Some(spec) = spec {
1287 return Source::Placeholder {
1288 what: format!(
1289 "datui read this through the format spec {spec}, which it cannot write as \
1290 Python: load it here."
1291 ),
1292 };
1293 }
1294 let targets: Option<Vec<Target>> = paths.iter().map(|p| reader_target(p, record)).collect();
1295 let Some(targets) = targets.filter(|t| !t.is_empty()) else {
1296 let names: Vec<String> = paths
1297 .iter()
1298 .map(|p| without_secrets(&p.to_string_lossy()).0)
1299 .collect();
1300 return Source::Placeholder {
1301 what: format!(
1302 "{}: datui could not name a Polars reader for this data; load it here.",
1303 named_with_table(&names, record)
1304 ),
1305 };
1306 };
1307 let format = targets[0].format;
1308 if targets.iter().any(|t| t.format != format) {
1309 return Source::Placeholder {
1310 what: "The files are of more than one format: load them here.".to_string(),
1311 };
1312 }
1313 let below = targets.iter().any(|t| t.below);
1314 let python = python_of(format);
1315 let literal = targets.iter().any(|t| t.literal);
1318 let no_glob = literal
1319 && python.is_some_and(|python| python.glob_flag)
1320 && !targets.iter().any(|t| t.pattern);
1321 let names: Vec<String> = targets
1322 .into_iter()
1323 .map(|t| {
1324 if t.literal && !no_glob {
1325 crate::cloud::source::escape_glob(&t.text)
1326 } else {
1327 t.text
1328 }
1329 })
1330 .collect();
1331 let Some(python) = python else {
1332 return Source::Placeholder {
1333 what: format!(
1334 "{}: Polars has no reader for {} files; load it here.",
1335 named_with_table(&names, record),
1336 format.title()
1337 ),
1338 };
1339 };
1340 let target = match names.as_slice() {
1341 [one] => py_str(one),
1342 many => py_names(many),
1343 };
1344 let options = record.options;
1345 let mut args: Vec<String> = vec![target];
1346 if no_glob {
1347 args.push("glob=False".to_string());
1348 }
1349 let mut notes = Vec::new();
1350 let endpoint = record.s3_endpoint.as_deref().map(without_secrets);
1351 if paths
1352 .iter()
1353 .any(|p| is_url(p) && without_secrets(&p.to_string_lossy()).1)
1354 || endpoint.as_ref().is_some_and(|(_, cut)| *cut)
1355 {
1356 notes.push(
1357 "datui left a user, password or query string out of the URL, as it may be a \
1358 credential: add it back if the server needs it."
1359 .to_string(),
1360 );
1361 }
1362 let in_store = names.iter().find(|n| store_scheme(n).is_some());
1363 let storage = in_store
1364 .and_then(|name| storage_options(name, record, endpoint.as_ref().map(|(e, _)| e.as_str())));
1365 if let Some(name) = in_store
1367 && python.eager
1368 {
1369 notes.push(format!(
1370 "{} reads no object store: download {name} and read it from disk.",
1371 python.call
1372 ));
1373 }
1374 let mut call = Call {
1375 record,
1376 paths,
1377 format,
1378 names: &names,
1379 below,
1380 storage,
1381 args,
1382 after: options.read_python.clone(),
1384 skip_tail: None,
1385 notes,
1386 };
1387 if let Some(arguments) = python.arguments
1388 && let Some(source) = arguments(&mut call)
1389 {
1390 return source;
1391 }
1392 let Call {
1393 args,
1394 mut after,
1395 skip_tail,
1396 mut notes,
1397 ..
1398 } = call;
1399 after.extend(skip_tail);
1400 if !record.read_as_text.is_empty() {
1401 notes.push(format!(
1402 "datui read these columns as text from every file: {}.",
1403 record.read_as_text.join(", ")
1404 ));
1405 }
1406 let mut call = format!("{}({})", python.call, args.join(", "));
1407 if python.eager {
1408 notes.extend(read_whole_note(record, &call));
1409 call.push_str(".lazy()");
1410 }
1411 Source::Read {
1412 call,
1413 after,
1414 notes,
1415 imports: Vec::new(),
1416 }
1417}
1418
1419fn csv_null_values(options: &OpenOptions, schema: &Schema) -> Option<String> {
1423 let specs = options.null_values.as_ref().filter(|s| !s.is_empty())?;
1424 let mut global = Vec::new();
1425 let mut per_column: Vec<(String, String)> = Vec::new();
1426 for spec in specs {
1427 match spec.find('=') {
1428 Some(i) => per_column.push((spec[..i].to_string(), spec[i + 1..].to_string())),
1429 None => global.push(spec.clone()),
1430 }
1431 }
1432 let dict = |pairs: Vec<(String, String)>| {
1433 let items: Vec<String> = pairs
1434 .iter()
1435 .map(|(c, v)| format!("{}: {}", py_str(c), py_str(v)))
1436 .collect();
1437 format!("{{{}}}", items.join(", "))
1438 };
1439 Some(match (global.as_slice(), per_column.is_empty()) {
1440 ([one], true) => py_str(one),
1441 (_, true) => py_names(&global),
1442 ([], false) => dict(per_column),
1443 (_, false) => dict(
1444 schema
1445 .iter_names()
1446 .map(|name| {
1447 let value = per_column
1448 .iter()
1449 .rev()
1450 .find(|(c, _)| c == name.as_str())
1451 .map(|(_, v)| v.clone())
1452 .unwrap_or_else(|| global[0].clone());
1453 (name.to_string(), value)
1454 })
1455 .collect(),
1456 ),
1457 })
1458}
1459
1460pub fn py_value(value: &AnyValue) -> Option<String> {
1463 Some(match value {
1464 AnyValue::Null => "None".to_string(),
1465 AnyValue::Boolean(b) => py_bool(*b).to_string(),
1466 AnyValue::String(s) => py_str(s),
1467 AnyValue::StringOwned(s) => py_str(s),
1468 AnyValue::Int8(v) => v.to_string(),
1469 AnyValue::Int16(v) => v.to_string(),
1470 AnyValue::Int32(v) => v.to_string(),
1471 AnyValue::Int64(v) => v.to_string(),
1472 AnyValue::UInt8(v) => v.to_string(),
1473 AnyValue::UInt16(v) => v.to_string(),
1474 AnyValue::UInt32(v) => v.to_string(),
1475 AnyValue::UInt64(v) => v.to_string(),
1476 AnyValue::Float32(v) => py_float(f64::from(*v)),
1477 AnyValue::Float64(v) => py_float(*v),
1478 AnyValue::Date(days) => {
1479 let date = chrono::NaiveDate::from_ymd_opt(1970, 1, 1)?
1480 .checked_add_signed(chrono::Duration::days(i64::from(*days)))?;
1481 use chrono::Datelike;
1482 format!("pl.date({}, {}, {})", date.year(), date.month(), date.day())
1483 }
1484 _ => return None,
1485 })
1486}
1487
1488#[derive(Debug, Clone, PartialEq)]
1490pub struct Script {
1491 pub source: Source,
1492 pub steps: Vec<Step>,
1493}
1494
1495impl Script {
1496 pub fn render(&self) -> String {
1497 let mut out = String::from("import polars as pl\n");
1498 if let Source::Read { imports, .. } = &self.source {
1499 for import in imports {
1500 out.push_str(import);
1501 out.push('\n');
1502 }
1503 }
1504 out.push('\n');
1505 let (head, mut lines) = match &self.source {
1506 Source::Read {
1507 call, after, notes, ..
1508 } => {
1509 for note in notes {
1510 out.push_str(&py_comment(note));
1511 out.push('\n');
1512 }
1513 (call.clone(), after.clone())
1514 }
1515 Source::Placeholder { what } => {
1516 out.push_str(&py_comment(what));
1517 out.push_str("\ndf = ...\n\n");
1518 ("df.lazy()".to_string(), Vec::new())
1519 }
1520 };
1521 let mut stopped = false;
1524 for step in &self.steps {
1525 let calls = step.python();
1526 if stopped {
1527 for call in &calls {
1529 lines.extend(call.lines().map(|c| {
1530 if c.starts_with('#') {
1531 c.to_string()
1532 } else {
1533 format!("# {c}")
1534 }
1535 }));
1536 }
1537 } else {
1538 stopped = matches!(step, Step::Unreproducible(_));
1539 lines.extend(calls);
1540 }
1541 }
1542 if lines.is_empty() {
1543 out.push_str(&format!("df = {head}\n"));
1544 } else {
1545 out.push_str("df = (\n");
1546 out.push_str(&format!(" {head}\n"));
1547 for line in lines {
1548 out.push_str(&format!(" {line}\n"));
1549 }
1550 out.push_str(")\n");
1551 }
1552 out
1553 }
1554}
1555
1556#[cfg(test)]
1557mod tests;