1use super::*;
5
6pub type FileScan = Arc<dyn Fn(&[String], &[PlSmallStr]) -> PolarsResult<LazyFrame> + Send + Sync>;
9
10pub type FileCounter = Arc<
12 dyn Fn(&Arc<crate::formats::schema_union::FooterProgress>) -> Result<Vec<Vec<usize>>, String>
13 + Send
14 + Sync,
15>;
16
17pub type FootersJoin = Arc<
20 dyn Fn(&Arc<crate::formats::schema_union::FooterProgress>) -> Option<FootersFound>
21 + Send
22 + Sync,
23>;
24
25pub struct FootersFound {
28 pub dataset: crate::formats::schema_union::DatasetSchema,
30 pub lf: LazyFrame,
32 pub file_rows: Vec<usize>,
34 pub files: Vec<String>,
37 pub row_groups: Vec<Vec<usize>>,
39 pub remote: Option<RemoteRead>,
41 pub estimate: Option<crate::formats::schema_union::RowEstimate>,
43}
44
45pub struct RemoteRead {
50 pub urls: Vec<String>,
52 pub scan: FileScan,
53 pub count: FileCounter,
54}
55
56impl From<RemoteRead> for RemoteFiles {
57 fn from(read: RemoteRead) -> Self {
58 RemoteFiles {
59 urls: Arc::new(read.urls),
60 scan: read.scan,
61 count: read.count,
62 offsets: None,
63 }
64 }
65}
66
67#[derive(Clone)]
72pub struct RemoteFiles {
73 pub urls: Arc<Vec<String>>,
75 pub scan: FileScan,
76 pub count: FileCounter,
77 pub offsets: Option<Vec<usize>>,
79}
80
81pub(super) fn conflicting_row_runs(
89 starts: &[usize],
90 total: usize,
91 conflicts: &[bool],
92) -> Vec<(usize, usize)> {
93 let mut runs: Vec<(usize, usize)> = Vec::new();
94 for (file, start) in starts.iter().copied().enumerate() {
95 if !conflicts.get(file).copied().unwrap_or(false) {
96 continue;
97 }
98 let end = starts.get(file + 1).copied().unwrap_or(total);
99 match runs.last_mut() {
100 Some(last) if last.1 == start => last.1 = end,
101 _ => runs.push((start, end)),
102 }
103 }
104 runs.retain(|(start, end)| start < end);
106 runs
107}
108
109impl DataTableState {
110 pub(crate) fn invalidate_num_rows(&mut self) {
113 self.view.num_rows_valid = false;
114 self.view.len_generation = next_len_generation();
115 }
116
117 pub(crate) fn changes_rows(&self) -> bool {
120 !self.view.filters.is_empty()
121 || !self.view.active_query.is_empty()
122 || !self.view.active_sql_query.is_empty()
123 || !self.view.active_fuzzy_query.is_empty()
124 || self.view.reshaped_lf.is_some()
125 || self.view.grouped.is_some()
126 || self.view.drilled_down_group_index.is_some()
127 }
128
129 pub(crate) fn may_keep_scan_rows(&self) -> bool {
132 self.view.filters.is_empty()
133 && self.view.active_fuzzy_query.is_empty()
134 && self.view.reshaped_lf.is_none()
135 && self.view.grouped.is_none()
136 && self.view.drilled_down_group_index.is_none()
137 }
138
139 pub(super) fn is_pristine(&self) -> bool {
142 self.view.column_changes.is_empty()
143 && self.view.filters.is_empty()
144 && self.view.sort_columns.is_empty()
145 && self.view.sort_ascending
146 && self.view.active_query.is_empty()
147 && self.view.active_sql_query.is_empty()
148 && self.view.active_fuzzy_query.is_empty()
149 && self.view.reshaped_lf.is_none()
150 && self.view.grouped.is_none()
151 && self.view.drilled_down_group_index.is_none()
152 }
153
154 pub fn analysis_lf(&self) -> LazyFrame {
157 self.view
158 .unsorted_lf
159 .clone()
160 .unwrap_or_else(|| self.view.lf.clone())
161 }
162
163 pub fn preview_lf(&self) -> LazyFrame {
166 Self::without_drift(self.analysis_lf())
167 }
168
169 pub fn is_sorted(&self) -> bool {
171 self.view.unsorted_lf.is_some()
172 }
173
174 pub fn scan_is_the_root(&self) -> bool {
179 self.view.active_query.is_empty()
180 && self.view.active_sql_query.is_empty()
181 && self.view.active_fuzzy_query.is_empty()
182 && self.view.reshaped_lf.is_none()
183 && self.view.grouped.is_none()
184 && self.view.drilled_down_group_index.is_none()
185 }
186
187 pub(super) fn restore_footer_count(&mut self) {
190 if !self.is_pristine() {
191 return;
192 }
193 if let Some(total) = self.row_group_offsets.as_ref().and_then(|o| o.last()) {
194 self.set_num_rows(*total);
195 }
196 }
197
198 pub fn measurements(&self) -> &Arc<crate::loading::measurements::Meter> {
200 &self.measurements
201 }
202
203 pub fn parquet_count_dir(&self) -> Option<PathBuf> {
206 self.parquet_count_dir
207 .clone()
208 .filter(|_| self.is_pristine())
209 }
210
211 pub fn len_generation(&self) -> u64 {
214 self.view.len_generation
215 }
216
217 pub fn num_rows_if_valid(&self) -> Option<usize> {
219 if self.view.num_rows_valid {
220 Some(self.view.num_rows)
221 } else {
222 None
223 }
224 }
225
226 pub fn is_num_rows_valid(&self) -> bool {
229 self.view.num_rows_valid
230 }
231
232 pub(super) fn num_rows_bound(&self) -> usize {
235 if self.view.num_rows_valid {
236 self.view.num_rows
237 } else {
238 usize::MAX
239 }
240 }
241
242 pub(super) fn set_num_rows(&mut self, n: usize) {
244 self.view.num_rows = n;
245 self.view.num_rows_valid = true;
246 self.remember_pristine_count();
247 if self.view.start_row > 0 && self.view.start_row >= n {
249 self.view.start_row = n.saturating_sub(self.visible_rows);
250 self.needs_recollect = true;
251 }
252 }
253
254 pub(super) fn remember_pristine_count(&mut self) {
257 if self.view.num_rows_valid && self.error.is_none() && self.is_pristine() {
258 self.pristine_rows = Some(self.view.num_rows);
259 }
260 }
261
262 pub fn lf_clone(&self) -> LazyFrame {
264 self.view.lf.clone()
265 }
266
267 pub fn polars_streaming_enabled(&self) -> bool {
269 self.polars_streaming
270 }
271
272 pub(crate) fn stitches_buffer(&self) -> bool {
276 self.is_pristine() && self.buffer_on_hand() && self.buffer_has_the_columns()
277 }
278
279 fn buffer_has_the_columns(&self) -> bool {
282 self.view.buffered_df.as_ref().is_some_and(|df| {
283 let names = df.columns().iter().map(|c| c.name().as_str());
284 let mut shown = self.view.column_order.iter().map(String::as_str);
285 names
286 .filter(|name| *name != crate::formats::schema_union::DRIFT_COLUMN)
287 .all(|name| shown.next() == Some(name))
288 && shown.next().is_none()
289 })
290 }
291
292 pub(crate) fn rows_on_hand(&self) -> Option<(&DataFrame, usize)> {
295 self.view
296 .buffered_df
297 .as_ref()
298 .filter(|_| self.buffer_on_hand())
299 .map(|df| (df, self.view.buffered_start_row))
300 }
301
302 pub(crate) fn buffer_on_hand(&self) -> bool {
304 self.view.buffered_end_row > self.view.buffered_start_row
305 && self.view.buffered_df.as_ref().is_some_and(|b| {
306 b.height() == self.view.buffered_end_row - self.view.buffered_start_row
307 })
308 }
309
310 pub(super) fn holds_buffer(&mut self, start: usize, end: usize) -> bool {
313 if !self.buffer_on_hand()
314 || start < self.view.buffered_start_row
315 || end > self.view.buffered_end_row
316 || end <= start
317 {
318 return false;
319 }
320 if (start, end) != (self.view.buffered_start_row, self.view.buffered_end_row) {
321 let offset = start - self.view.buffered_start_row;
322 self.view.locked_df = None;
325 self.view.df = None;
326 self.view.buffered_df = self
327 .view
328 .buffered_df
329 .take()
330 .map(|b| trim_rows(b, offset, end - start, None));
331 self.view.buffered_start_row = start;
332 self.view.buffered_end_row = end;
333 }
334 true
335 }
336
337 pub fn buffered_start(&self) -> usize {
339 self.view.buffered_start_row
340 }
341
342 pub fn buffered_end(&self) -> usize {
344 self.view.buffered_end_row
345 }
346
347 pub fn is_remote_source(&self) -> bool {
352 self.remote_source
353 }
354
355 pub fn source_schema(&self) -> &Arc<Schema> {
358 &self.original_schema
359 }
360
361 pub(super) fn record_row_groups(&mut self, rows: &[usize]) {
364 let mut offsets = Vec::with_capacity(rows.len() + 1);
365 offsets.push(0);
366 for n in rows {
367 offsets.push(offsets.last().unwrap_or(&0) + n);
368 }
369 self.set_num_rows(*offsets.last().unwrap_or(&0));
370 self.row_group_offsets = Some(offsets);
371 }
372
373 pub(crate) fn quality_reads_whole_source(
377 &self,
378 scope: &crate::analysis::data_quality::QualityScope,
379 ) -> bool {
380 use crate::analysis::data_quality::QualityScope;
381 let columns = || {
382 self.original_schema
383 .iter()
384 .filter(|(name, _)| name.as_str() != crate::formats::schema_union::DRIFT_COLUMN)
385 };
386 if columns().any(|(_, dtype)| matches!(dtype, DataType::Binary)) {
387 return false;
388 }
389 match scope {
390 QualityScope::WholeSource => true,
391 QualityScope::CurrentView => {
392 let shown = self
393 .view
394 .column_order
395 .iter()
396 .map(String::as_str)
397 .collect::<HashSet<_>>();
398 !self.changes_rows() && columns().all(|(name, _)| shown.contains(name.as_str()))
399 }
400 _ => false,
401 }
402 }
403
404 pub(crate) fn each_remote_object(
407 &self,
408 ) -> Option<Box<dyn Iterator<Item = Option<&RemoteObject>> + '_>> {
409 let objects = self.remote_objects.as_ref()?;
410 Some(match &self.remote_files {
411 Some(remote) => Box::new(remote.urls.iter().map(|url| objects.get(url))),
412 None => Box::new(objects.values().map(Some)),
413 })
414 }
415
416 pub(crate) fn remote_objects(&self) -> Option<Vec<RemoteObject>> {
418 let objects = self
419 .each_remote_object()?
420 .map(|object| object.cloned())
421 .collect::<Option<Vec<_>>>()?;
422 (!objects.is_empty()).then_some(objects)
423 }
424
425 pub(crate) fn remote_objects_size(&self) -> Option<(u64, usize)> {
427 let (bytes, count) = self
428 .each_remote_object()?
429 .try_fold((0u64, 0usize), |(bytes, count), object| {
430 object.map(|object| (bytes + object.size, count + 1))
431 })?;
432 (count > 0).then_some((bytes, count))
433 }
434
435 pub(super) fn record_dataset_schema(
439 &mut self,
440 schema: crate::formats::schema_union::DatasetSchema,
441 file_rows: &[usize],
442 files: &[String],
443 ) {
444 self.drift_files = files.to_vec();
445 self.view.drift_column_present =
447 schema.drifts() && file_rows.len() == schema.file_group.len();
448 self.view.drift_groups = Arc::new(schema.groups.clone());
449 self.drift_file_group = schema.file_group.clone();
450 self.drift_file_starts = Vec::with_capacity(file_rows.len());
451 let mut row = 0usize;
452 for rows in file_rows {
453 self.drift_file_starts.push(row);
454 row += rows;
455 }
456 self.drift_dataset_rows = row;
457 self.drift_at_open = self.view.drift_column_present;
458 self.groups_at_open = self.view.drift_groups.clone();
459 self.view.notes = Self::notes_datui_can_act_on(&schema, self.view.drift_column_present);
460 self.notes_at_open = self.view.notes.clone();
461 self.view.notes_seen = false;
462 self.read_as_text = Vec::new();
463 self.dataset_at_open = Some(schema.clone());
464 self.dataset_schema = Some(schema);
465 }
466
467 pub fn footers_pending(&self) -> Option<FootersJoin> {
469 self.footers_pending.clone()
470 }
471
472 pub fn counts_itself_later(&self) -> bool {
476 if self.indexing().is_some() && !self.view.num_rows_valid {
479 return true;
480 }
481 self.footers_pending.is_some() && !self.view.num_rows_valid && self.is_pristine()
484 }
485
486 pub fn indexing(&self) -> Option<&Arc<crate::formats::lines::Lines>> {
488 self.indexing.as_ref().filter(|lines| lines.indexing())
490 }
491
492 pub fn lines_to_index(&self) -> Option<&Arc<crate::formats::lines::Lines>> {
495 self.indexing.as_ref()
496 }
497
498 pub fn row_estimate(
501 &self,
502 pass: Option<crate::formats::schema_union::RowEstimate>,
503 ) -> Option<crate::formats::schema_union::RowEstimate> {
504 if self.view.num_rows_valid || !self.is_pristine() {
505 return None;
506 }
507 self.row_estimate
508 .or_else(|| pass.filter(|_| self.footers_pending.is_some()))
509 }
510
511 pub fn file_row_starts(&self) -> Option<Vec<usize>> {
514 if self.changes_rows() {
515 return None;
516 }
517 self.remote_files.as_ref()?.offsets.clone()
518 }
519
520 pub fn files_to_count(&self) -> Option<usize> {
522 self.remote_files
523 .as_ref()
524 .filter(|f| f.offsets.is_none())
525 .map(|f| f.urls.len())
526 }
527
528 pub fn numbered_by_default(&self) -> bool {
531 matches!(
532 self.read_as,
533 Some(crate::FileFormat::Text | crate::FileFormat::Journal)
534 )
535 }
536
537 pub fn row_numbers_count_the_view(&self) -> bool {
540 self.row_numbers
541 && !self.carries_source_rows()
542 && self.scan_is_the_root()
543 && (!self.view.filters.is_empty()
544 || !self.view.sort_columns.is_empty()
545 || !self.view.sort_ascending)
546 }
547
548 pub(crate) fn lines_indexed(&mut self, rows: usize) -> bool {
552 let Some(lines) = self.indexing.take() else {
553 return false;
554 };
555 let notes = crate::formats::lines::notes(&lines, self.indexing_guessed);
556 let opened = std::mem::take(&mut self.indexing_notes);
557 self.open_notes.retain(|n| !opened.contains(n));
558 self.open_notes.extend(notes);
559 if lines.shrank() {
561 self.open_notes.push(crate::formats::text_formats::note(
562 crate::formats::lines::SHRANK.to_string(),
563 "the file".to_string(),
564 ));
565 return true;
566 }
567 self.pristine_rows = Some(rows);
569 if self.is_pristine() {
570 self.set_num_rows(rows);
571 }
572 true
573 }
574
575 pub fn give_up_on_pending_footers(&mut self) {
579 self.footers_pending = None;
580 }
581
582 pub fn join_dataset_schema(
589 &mut self,
590 mut found: FootersFound,
591 ) -> std::result::Result<(), Box<FootersFound>> {
592 if !self.scan_is_the_root() {
593 let row_groups = std::mem::take(&mut found.row_groups);
597 if !row_groups.is_empty() {
598 self.record_file_row_groups(&row_groups);
599 }
600 return Err(Box::new(found));
602 }
603 let FootersFound {
604 dataset,
605 lf,
606 file_rows,
607 files,
608 row_groups,
609 remote,
610 estimate,
611 } = found;
612 self.row_estimate = if row_groups.is_empty() {
613 estimate
614 } else {
615 None
616 };
617 let (file_rows, files) = (file_rows.as_slice(), files.as_slice());
618 let known: std::collections::HashSet<&str> =
619 self.view.column_order.iter().map(String::as_str).collect();
620 let joining: Vec<String> = dataset
621 .schema
622 .iter_names()
623 .map(|name| name.to_string())
624 .filter(|name| {
625 name != crate::formats::schema_union::DRIFT_COLUMN && !known.contains(name.as_str())
626 })
627 .collect();
628 drop(known);
629 self.view.column_order.extend(joining);
630 self.view
633 .column_order
634 .retain(|name| dataset.schema.contains(name.as_str()));
635 let schema = dataset.schema.clone();
636 match (remote, self.remote_files.as_mut()) {
639 (Some(found), Some(remote)) => {
640 remote.urls = Arc::new(found.urls);
641 remote.scan = found.scan;
642 remote.count = found.count;
645 }
646 (Some(found), None) => self.remote_files = Some(found.into()),
648 (None, _) => {}
649 }
650 self.record_dataset_schema(dataset, file_rows, files);
653 self.footers_pending = None;
654 self.replace_root(lf, schema);
657 self.pristine_rows = None;
659 self.view.observed_bytes_per_row = None;
662 if !row_groups.is_empty() {
665 self.record_file_row_groups(&row_groups);
666 }
667 self.deferred(Self::apply_transformations);
671 Ok(())
672 }
673
674 pub fn visible_lf(&self) -> LazyFrame {
678 Self::without_drift(self.view.lf.clone())
679 }
680
681 pub fn scans_a_temp_file(&self) -> bool {
684 self.decompress_temp_file.is_some() || !self.converted.is_empty()
685 }
686
687 pub fn scans_a_download(&self) -> bool {
690 self.download.is_some()
691 }
692
693 pub fn read_mode(&self) -> Option<crate::ReadMode> {
695 self.read_mode
696 }
697
698 pub fn read_as(&self) -> Option<crate::FileFormat> {
700 self.read_as
701 }
702
703 pub fn fetched(&self) -> bool {
705 self.fetched
706 }
707
708 pub(crate) fn temp_files(&self) -> Vec<&Path> {
711 let files = self.decompress_temp_file.iter().map(|file| file.path());
712 let files = files.chain(self.download.iter().map(|download| download.path()));
713 let files = files.chain(self.converted.iter().map(|file| file.path()));
714 files.collect()
715 }
716
717 pub(super) fn without_drift(lf: LazyFrame) -> LazyFrame {
719 lf.drop(by_name(
720 [crate::formats::schema_union::DRIFT_COLUMN],
721 false,
722 false,
723 ))
724 }
725
726 pub(super) fn query_source(&self) -> LazyFrame {
729 Self::without_drift(self.original_lf.clone())
730 }
731
732 pub(super) fn query_source_schema(&self) -> Arc<Schema> {
734 Self::without_source_rows(self.original_schema.clone()).0
735 }
736
737 pub fn drifts(&self) -> bool {
739 self.view.drift_column_present
740 }
741
742 pub fn drift_groups(&self) -> Arc<Vec<crate::formats::schema_union::DriftGroup>> {
744 self.view.drift_groups.clone()
745 }
746
747 pub fn can_name_source_files(&self) -> bool {
750 self.view.drift_column_present
751 && !self.drift_files.is_empty()
752 && self.drift_files.len() == self.drift_file_starts.len()
753 }
754
755 pub fn export_frame(&self, name_files: bool) -> ExportFrame {
759 if name_files && self.can_name_source_files() {
760 ExportFrame {
761 lf: self.view.lf.clone(),
762 files: Some(SourceFiles {
763 names: Arc::new(self.drift_files.clone()),
764 starts: Arc::new(self.drift_file_starts.clone()),
765 }),
766 }
767 } else {
768 ExportFrame {
769 lf: self.visible_lf(),
770 files: None,
771 }
772 }
773 }
774
775 pub fn notes(&self) -> Vec<crate::notes::Note> {
778 let mut notes = crate::notes::merged(
781 &self.open_notes,
782 &self.view.notes,
783 &self.view.view_notes,
784 self.dataset_schema.as_ref(),
785 );
786 if let Some(pushdown) = &self.pushdown {
788 notes.extend(pushdown.notes());
789 }
790 notes.extend(self.unfit_notes.iter().flatten().cloned());
791 notes.extend(self.view.changes_dropped.iter().cloned());
792 if let Some((version, unfit)) = &self.changes_unfit
793 && *version == self.view.changes_version
794 {
795 notes.extend(unfit.iter().cloned());
796 }
797 notes
798 }
799
800 pub(crate) fn window_for_quality(
803 &self,
804 scope: &crate::analysis::data_quality::QualityScope,
805 ) -> Option<Arc<dyn crate::formats::pushdown::Windowed>> {
806 use crate::analysis::data_quality::QualityScope;
807 matches!(scope, QualityScope::WholeSource | QualityScope::CurrentView)
808 .then(|| {
809 self.fixed_window
810 .clone()
811 .filter(|_| self.is_pristine() && self.indexing().is_none())
812 })
813 .flatten()
814 }
815
816 pub fn not_the_table(&self) -> Option<&'static str> {
819 self.not_the_table
820 }
821
822 pub fn other_tables(&self) -> &[String] {
824 &self.other_tables
825 }
826
827 pub fn format_read(&self) -> Option<&Arc<crate::formats::Read>> {
829 self.format_read.as_ref()
830 }
831
832 pub(super) fn window_now(&self) -> Option<Arc<dyn crate::formats::pushdown::Windowed>> {
835 if let Some(window) = self.follow_window() {
836 return Some(Arc::new(window));
837 }
838 if let Some(records) = self.fixed_window.as_ref().filter(|_| self.is_pristine()) {
839 return Some(records.clone());
840 }
841 self.pushed_view().map(|view| view.window)
842 }
843
844 pub(crate) fn pushed_view(&self) -> Option<crate::formats::pushdown::PushedView> {
847 let pushdown = self.pushdown.as_ref()?;
848 if !self.scan_is_the_root() || self.view.drift_column_present {
849 return None;
850 }
851 let sort: Vec<(String, bool)> = self
852 .view
853 .sort_columns
854 .iter()
855 .cloned()
856 .zip(self.view.sort_descending.iter().copied())
857 .collect();
858 pushdown.view(&self.view.filters, &sort, !self.view.sort_ascending)
859 }
860
861 pub(crate) fn source_counter(&self) -> Option<crate::formats::pushdown::Counter> {
864 if let Some(counter) = self.follow_counter() {
865 return Some(counter);
866 }
867 self.pushed_view().map(|view| view.counter)
868 }
869
870 fn follow_known(&self) -> Option<&[(usize, usize)]> {
873 self.follow_known
874 .as_ref()
875 .filter(|(generation, _)| *generation == self.view.len_generation)
876 .map(|(_, known)| known.as_slice())
877 }
878
879 fn follow_window(&self) -> Option<crate::loading::follow::Window> {
882 let follow = self.follow.as_ref()?;
883 let known = if self.is_pristine() {
884 None
885 } else if self.scan_is_the_root()
886 && self.view.sort_columns.is_empty()
887 && self.view.sort_ascending
888 {
889 Some(self.follow_known()?.to_vec())
890 } else {
891 return None;
892 };
893 Some(crate::loading::follow::Window {
894 lf: self.view.lf.clone(),
895 path: follow.path().to_path_buf(),
896 marks: follow.marks().clone(),
897 known,
898 })
899 }
900
901 fn follow_counter(&self) -> Option<crate::formats::pushdown::Counter> {
904 let follow = self.follow.as_ref()?;
905 let &(before, row) = self.follow_known()?.last()?;
906 let rest = crate::loading::follow::from_marks(
907 &self.view.lf,
908 follow.path(),
909 follow.marks(),
910 row,
911 None,
912 )?;
913 let streaming = self.polars_streaming;
914 Some(Arc::new(move || {
915 let df = crate::analysis::statistics::collect_lazy(row_count_lf(&rest), streaming)?;
916 let after = match df.get(0).and_then(|row| row.first().cloned()) {
917 Some(AnyValue::UInt64(n)) => n as usize,
918 _ => 0,
919 };
920 Ok(before + after)
921 }))
922 }
923
924 pub fn delimited_read(&self) -> Option<&Arc<crate::formats::delimited_spec::DelimitedRead>> {
927 self.delimited.as_ref()
928 }
929
930 pub fn unit_of(&self, column: &str) -> Option<&str> {
934 if self.delimited.is_none() && self.file_units.is_empty() {
935 return None;
936 }
937 let loaded = match &self.view.lineage {
938 None => column,
939 Some(lineage) => lineage
940 .iter()
941 .find(|(shown, _)| shown == column)
942 .map(|(_, loaded)| loaded.as_str())?,
943 };
944 match &self.delimited {
945 Some(read) => read.unit_of(loaded),
946 None => self
947 .file_units
948 .iter()
949 .find(|(name, _)| name == loaded)
950 .map(|(_, unit)| unit.as_str()),
951 }
952 }
953
954 pub fn units(&self) -> Vec<(String, String)> {
956 if self.delimited.is_none() && self.file_units.is_empty() {
957 return Vec::new();
958 }
959 self.view
960 .schema
961 .iter_names()
962 .filter_map(|name| Some((name.to_string(), self.unit_of(name)?.to_string())))
963 .collect()
964 }
965
966 pub fn format_detail(&self) -> Option<&crate::formats::text_formats::Detail> {
968 self.detail.as_deref()
969 }
970
971 pub(crate) fn ended_journal_to_describe(&mut self) -> Option<LazyFrame> {
974 let follow = self.follow.as_mut()?;
975 if follow.described
976 || follow.live()
977 || follow.behind()
978 || follow.spool().is_none()
979 || self.read_as != Some(crate::FileFormat::Journal)
980 {
981 return None;
982 }
983 follow.described = true;
984 Some(self.original_lf.clone())
985 }
986
987 pub(crate) fn set_format_detail(&mut self, detail: crate::formats::text_formats::Detail) {
988 self.detail = Some(Arc::new(detail));
989 }
990
991 pub fn has_notes(&self) -> bool {
993 !self.view.notes.is_empty()
996 || !self.open_notes.is_empty()
997 || self.unfit_notes.as_ref().is_some_and(|n| !n.is_empty())
998 || !self.view.changes_dropped.is_empty()
999 || self
1000 .changes_unfit
1001 .as_ref()
1002 .is_some_and(|(v, n)| *v == self.view.changes_version && !n.is_empty())
1003 || self
1004 .pushdown
1005 .as_ref()
1006 .is_some_and(|p| !p.notes().is_empty())
1007 }
1008
1009 pub fn notes_unseen(&self) -> bool {
1011 self.has_notes() && !self.view.notes_seen
1012 }
1013
1014 fn unread_row_runs(&self, column: &str) -> Vec<(usize, usize)> {
1017 let Some(dataset) = self.dataset_schema.as_ref() else {
1018 return Vec::new();
1019 };
1020 let conflicts: Vec<bool> = (0..self.drift_file_starts.len())
1021 .map(|file| {
1022 dataset
1023 .file_group
1024 .get(file)
1025 .and_then(|group| dataset.groups.get(*group as usize))
1026 .is_some_and(|group| group.unread.iter().any(|name| name == column))
1027 })
1028 .collect();
1029 conflicting_row_runs(&self.drift_file_starts, self.drift_dataset_rows, &conflicts)
1030 }
1031
1032 fn view_columns_with_conflicts(&self) -> Vec<crate::formats::schema_union::ColumnDrift> {
1035 let Some(dataset) = self.dataset_schema.as_ref() else {
1036 return Vec::new();
1037 };
1038 let named: HashSet<&str> = self
1039 .view
1040 .filters
1041 .iter()
1042 .map(|filter| filter.column.as_str())
1043 .chain(self.view.sort_columns.iter().map(String::as_str))
1044 .collect();
1045 dataset
1046 .columns
1047 .iter()
1048 .filter(|column| column.conflicting_files > 0 && named.contains(column.name.as_str()))
1049 .cloned()
1050 .collect()
1051 }
1052
1053 pub(super) fn view_exclusions(&self) -> Vec<(Vec<(usize, usize)>, crate::notes::Note)> {
1057 if !self.view.drift_column_present {
1058 return Vec::new();
1059 }
1060 let Some(dataset) = self.dataset_schema.as_ref() else {
1061 return Vec::new();
1062 };
1063 let mut out = Vec::new();
1064 for column in self.view_columns_with_conflicts() {
1065 let runs = self.unread_row_runs(&column.name);
1066 let rows: usize = runs.iter().map(|(start, end)| end - start).sum();
1067 if rows == 0 {
1068 continue;
1069 }
1070 let filtered = self
1071 .view
1072 .filters
1073 .iter()
1074 .any(|filter| filter.column.as_str() == column.name.as_str());
1075 let sorted = self
1076 .view
1077 .sort_columns
1078 .iter()
1079 .any(|sorted| sorted.as_str() == column.name.as_str());
1080 out.push((
1081 runs,
1082 crate::notes::left_out_note(&column, dataset, rows, filtered, sorted),
1083 ));
1084 }
1085 out
1086 }
1087
1088 pub(super) fn view_notes_only(&self) -> Vec<crate::notes::Note> {
1091 self.view_exclusions()
1092 .into_iter()
1093 .map(|(_, note)| note)
1094 .collect()
1095 }
1096
1097 pub(super) fn leave_out_unread_rows(
1098 &self,
1099 mut lf: LazyFrame,
1100 ) -> (LazyFrame, Vec<crate::notes::Note>) {
1101 let mut notes = Vec::new();
1102 for (runs, note) in self.view_exclusions() {
1103 let keep = runs
1104 .iter()
1105 .map(|(start, end)| {
1106 col(crate::formats::schema_union::DRIFT_COLUMN)
1107 .lt(lit(*start as u32))
1108 .or(col(crate::formats::schema_union::DRIFT_COLUMN).gt_eq(lit(*end as u32)))
1109 })
1110 .reduce(Expr::and);
1111 if let Some(keep) = keep {
1112 lf = lf.filter(keep);
1113 }
1114 notes.push(note);
1115 }
1116 (lf, notes)
1117 }
1118
1119 pub fn mark_notes_seen(&mut self) {
1121 self.view.notes_seen = true;
1122 }
1123
1124 pub fn read_column_as_text(&mut self, column: &str) -> PolarsResult<bool> {
1130 let name = PlSmallStr::from(column);
1131 let Some(dataset) = self.dataset_at_open.clone() else {
1132 return Ok(false);
1133 };
1134 if !self.view.drift_column_present || self.read_as_text.contains(&name) {
1135 return Ok(false);
1136 }
1137 if !dataset
1138 .columns
1139 .iter()
1140 .any(|drift| drift.name == name && drift.can_read_as_text())
1141 {
1142 return Ok(false);
1143 }
1144
1145 let mut as_text = self.read_as_text.clone();
1146 as_text.push(name);
1147
1148 let drift = crate::formats::schema_union::ScanDrift::new(
1151 &self.drift_files,
1152 &dataset,
1153 &self.file_rows(),
1154 );
1155 let scanned = match self.remote_files.as_ref() {
1156 Some(remote) => (remote.scan)(&remote.urls, &as_text),
1157 None => crate::formats::schema_union::lenient_scan(
1158 &self.drift_files,
1159 dataset.schema.clone(),
1160 None,
1161 drift.as_ref(),
1162 &as_text,
1163 ),
1164 };
1165 let lf = match scanned {
1166 Ok(lf) => lf,
1167 Err(e) => {
1168 self.error = Some(e.clone());
1170 return Err(e);
1171 }
1172 };
1173 let lf = if self.remote_files.is_some() {
1176 lf
1177 } else {
1178 crate::loading::open_scan::hoist_partition_columns(
1179 lf,
1180 &dataset.schema,
1181 self.partition_columns.as_deref().unwrap_or(&[]),
1182 drift.is_some(),
1183 )
1184 };
1185
1186 let view = dataset.reading_as_text(&as_text);
1187 self.read_as_text = as_text;
1188 let schema = view.schema.clone();
1190 self.view.drift_groups = Arc::new(view.groups.clone());
1191 self.groups_at_open = self.view.drift_groups.clone();
1192 self.view.notes = Self::notes_datui_can_act_on(&view, self.view.drift_column_present);
1193 self.notes_at_open = self.view.notes.clone();
1194 self.view.notes_seen = false;
1196 self.dataset_schema = Some(view);
1197 self.replace_root(lf, schema);
1198 self.apply_transformations();
1200 Ok(true)
1201 }
1202
1203 fn notes_datui_can_act_on(
1207 dataset: &crate::formats::schema_union::DatasetSchema,
1208 counted: bool,
1209 ) -> Vec<crate::notes::Note> {
1210 let mut notes = crate::notes::from_dataset(dataset);
1211 if !counted {
1212 for note in &mut notes {
1213 note.read_as_text = None;
1214 }
1215 }
1216 notes
1217 }
1218
1219 pub(super) fn file_rows(&self) -> Vec<usize> {
1222 self.drift_file_starts
1223 .iter()
1224 .enumerate()
1225 .map(|(file, start)| {
1226 self.drift_file_starts
1227 .get(file + 1)
1228 .copied()
1229 .unwrap_or(self.drift_dataset_rows)
1230 .saturating_sub(*start)
1231 })
1232 .collect()
1233 }
1234
1235 pub fn read_as_text(&self) -> &[PlSmallStr] {
1237 &self.read_as_text
1238 }
1239
1240 pub fn dataset_schema(&self) -> Option<&crate::formats::schema_union::DatasetSchema> {
1242 self.dataset_schema.as_ref()
1243 }
1244
1245 pub fn remote_files_counter(&self) -> Option<FileCounter> {
1249 if self.footers_pending.is_some() {
1250 return None;
1251 }
1252 self.remote_files
1253 .as_ref()
1254 .filter(|f| f.offsets.is_none() && self.is_pristine())
1255 .map(|f| f.count.clone())
1256 }
1257
1258 pub(super) fn record_file_row_groups(&mut self, groups: &[Vec<usize>]) {
1261 if let Some(files) = self.remote_files.as_mut() {
1262 if groups.len() != files.urls.len() {
1263 return;
1264 }
1265 let mut offsets = Vec::with_capacity(groups.len() + 1);
1266 offsets.push(0);
1267 for file in groups {
1268 offsets.push(offsets.last().unwrap_or(&0) + file.iter().sum::<usize>());
1269 }
1270 files.offsets = Some(offsets);
1271 }
1272 let flat: Vec<usize> = groups.iter().flatten().copied().collect();
1273 if self.is_pristine() {
1274 self.record_row_groups(&flat);
1275 } else {
1276 let mut row_offsets = Vec::with_capacity(flat.len() + 1);
1278 row_offsets.push(0);
1279 for n in &flat {
1280 row_offsets.push(row_offsets.last().unwrap_or(&0) + n);
1281 }
1282 self.row_group_offsets = Some(row_offsets);
1283 }
1284 }
1285}