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 {
275 self.remote_window() && self.buffer_on_hand()
276 }
277
278 pub(crate) fn rows_on_hand(&self) -> Option<(&DataFrame, usize)> {
281 self.view
282 .buffered_df
283 .as_ref()
284 .filter(|_| self.buffer_on_hand())
285 .map(|df| (df, self.view.buffered_start_row))
286 }
287
288 pub(crate) fn buffer_on_hand(&self) -> bool {
290 self.view.buffered_end_row > self.view.buffered_start_row
291 && self.view.buffered_df.as_ref().is_some_and(|b| {
292 b.height() == self.view.buffered_end_row - self.view.buffered_start_row
293 })
294 }
295
296 pub(super) fn holds_buffer(&mut self, start: usize, end: usize) -> bool {
299 if !self.buffer_on_hand()
300 || start < self.view.buffered_start_row
301 || end > self.view.buffered_end_row
302 || end <= start
303 {
304 return false;
305 }
306 if (start, end) != (self.view.buffered_start_row, self.view.buffered_end_row) {
307 let offset = start - self.view.buffered_start_row;
308 self.view.locked_df = None;
311 self.view.df = None;
312 self.view.buffered_df = self
313 .view
314 .buffered_df
315 .take()
316 .map(|b| trim_rows(b, offset, end - start, None));
317 self.view.buffered_start_row = start;
318 self.view.buffered_end_row = end;
319 }
320 true
321 }
322
323 pub fn buffered_start(&self) -> usize {
325 self.view.buffered_start_row
326 }
327
328 pub fn buffered_end(&self) -> usize {
330 self.view.buffered_end_row
331 }
332
333 pub fn is_remote_source(&self) -> bool {
338 self.remote_source
339 }
340
341 pub fn source_schema(&self) -> &Arc<Schema> {
344 &self.original_schema
345 }
346
347 pub(super) fn record_row_groups(&mut self, rows: &[usize]) {
350 let mut offsets = Vec::with_capacity(rows.len() + 1);
351 offsets.push(0);
352 for n in rows {
353 offsets.push(offsets.last().unwrap_or(&0) + n);
354 }
355 self.set_num_rows(*offsets.last().unwrap_or(&0));
356 self.row_group_offsets = Some(offsets);
357 }
358
359 pub(crate) fn quality_reads_whole_source(
363 &self,
364 scope: &crate::analysis::data_quality::QualityScope,
365 ) -> bool {
366 use crate::analysis::data_quality::QualityScope;
367 let columns = || {
368 self.original_schema
369 .iter()
370 .filter(|(name, _)| name.as_str() != crate::formats::schema_union::DRIFT_COLUMN)
371 };
372 if columns().any(|(_, dtype)| matches!(dtype, DataType::Binary)) {
373 return false;
374 }
375 match scope {
376 QualityScope::WholeSource => true,
377 QualityScope::CurrentView => {
378 let shown = self
379 .view
380 .column_order
381 .iter()
382 .map(String::as_str)
383 .collect::<HashSet<_>>();
384 !self.changes_rows() && columns().all(|(name, _)| shown.contains(name.as_str()))
385 }
386 _ => false,
387 }
388 }
389
390 pub(crate) fn each_remote_object(
393 &self,
394 ) -> Option<Box<dyn Iterator<Item = Option<&RemoteObject>> + '_>> {
395 let objects = self.remote_objects.as_ref()?;
396 Some(match &self.remote_files {
397 Some(remote) => Box::new(remote.urls.iter().map(|url| objects.get(url))),
398 None => Box::new(objects.values().map(Some)),
399 })
400 }
401
402 pub(crate) fn remote_objects(&self) -> Option<Vec<RemoteObject>> {
404 let objects = self
405 .each_remote_object()?
406 .map(|object| object.cloned())
407 .collect::<Option<Vec<_>>>()?;
408 (!objects.is_empty()).then_some(objects)
409 }
410
411 pub(crate) fn remote_objects_size(&self) -> Option<(u64, usize)> {
413 let (bytes, count) = self
414 .each_remote_object()?
415 .try_fold((0u64, 0usize), |(bytes, count), object| {
416 object.map(|object| (bytes + object.size, count + 1))
417 })?;
418 (count > 0).then_some((bytes, count))
419 }
420
421 pub(super) fn record_dataset_schema(
425 &mut self,
426 schema: crate::formats::schema_union::DatasetSchema,
427 file_rows: &[usize],
428 files: &[String],
429 ) {
430 self.drift_files = files.to_vec();
431 self.view.drift_column_present =
433 schema.drifts() && file_rows.len() == schema.file_group.len();
434 self.view.drift_groups = Arc::new(schema.groups.clone());
435 self.drift_file_group = schema.file_group.clone();
436 self.drift_file_starts = Vec::with_capacity(file_rows.len());
437 let mut row = 0usize;
438 for rows in file_rows {
439 self.drift_file_starts.push(row);
440 row += rows;
441 }
442 self.drift_dataset_rows = row;
443 self.drift_at_open = self.view.drift_column_present;
444 self.groups_at_open = self.view.drift_groups.clone();
445 self.view.notes = Self::notes_datui_can_act_on(&schema, self.view.drift_column_present);
446 self.notes_at_open = self.view.notes.clone();
447 self.view.notes_seen = false;
448 self.read_as_text = Vec::new();
449 self.dataset_at_open = Some(schema.clone());
450 self.dataset_schema = Some(schema);
451 }
452
453 pub fn footers_pending(&self) -> Option<FootersJoin> {
455 self.footers_pending.clone()
456 }
457
458 pub fn counts_itself_later(&self) -> bool {
462 if self.indexing().is_some() && !self.view.num_rows_valid {
465 return true;
466 }
467 self.footers_pending.is_some() && !self.view.num_rows_valid && self.is_pristine()
470 }
471
472 pub fn indexing(&self) -> Option<&Arc<crate::formats::lines::Lines>> {
474 self.indexing.as_ref().filter(|lines| lines.indexing())
476 }
477
478 pub fn lines_to_index(&self) -> Option<&Arc<crate::formats::lines::Lines>> {
481 self.indexing.as_ref()
482 }
483
484 pub fn row_estimate(
487 &self,
488 pass: Option<crate::formats::schema_union::RowEstimate>,
489 ) -> Option<crate::formats::schema_union::RowEstimate> {
490 if self.view.num_rows_valid || !self.is_pristine() {
491 return None;
492 }
493 self.row_estimate
494 .or_else(|| pass.filter(|_| self.footers_pending.is_some()))
495 }
496
497 pub fn file_row_starts(&self) -> Option<Vec<usize>> {
500 if self.changes_rows() {
501 return None;
502 }
503 self.remote_files.as_ref()?.offsets.clone()
504 }
505
506 pub fn files_to_count(&self) -> Option<usize> {
508 self.remote_files
509 .as_ref()
510 .filter(|f| f.offsets.is_none())
511 .map(|f| f.urls.len())
512 }
513
514 pub fn numbered_by_default(&self) -> bool {
517 matches!(
518 self.read_as,
519 Some(crate::FileFormat::Text | crate::FileFormat::Journal)
520 )
521 }
522
523 pub fn row_numbers_count_the_view(&self) -> bool {
526 self.row_numbers
527 && !self.carries_source_rows()
528 && self.scan_is_the_root()
529 && (!self.view.filters.is_empty()
530 || !self.view.sort_columns.is_empty()
531 || !self.view.sort_ascending)
532 }
533
534 pub(crate) fn lines_indexed(&mut self, rows: usize) -> bool {
538 let Some(lines) = self.indexing.take() else {
539 return false;
540 };
541 let notes = crate::formats::lines::notes(&lines, self.indexing_guessed);
542 let opened = std::mem::take(&mut self.indexing_notes);
543 self.open_notes.retain(|n| !opened.contains(n));
544 self.open_notes.extend(notes);
545 if lines.shrank() {
547 self.open_notes.push(crate::formats::text_formats::note(
548 crate::formats::lines::SHRANK.to_string(),
549 "the file".to_string(),
550 ));
551 return true;
552 }
553 self.pristine_rows = Some(rows);
555 if self.is_pristine() {
556 self.set_num_rows(rows);
557 }
558 true
559 }
560
561 pub fn give_up_on_pending_footers(&mut self) {
565 self.footers_pending = None;
566 }
567
568 pub fn join_dataset_schema(
575 &mut self,
576 mut found: FootersFound,
577 ) -> std::result::Result<(), Box<FootersFound>> {
578 if !self.scan_is_the_root() {
579 let row_groups = std::mem::take(&mut found.row_groups);
583 if !row_groups.is_empty() {
584 self.record_file_row_groups(&row_groups);
585 }
586 return Err(Box::new(found));
588 }
589 let FootersFound {
590 dataset,
591 lf,
592 file_rows,
593 files,
594 row_groups,
595 remote,
596 estimate,
597 } = found;
598 self.row_estimate = if row_groups.is_empty() {
599 estimate
600 } else {
601 None
602 };
603 let (file_rows, files) = (file_rows.as_slice(), files.as_slice());
604 let known: std::collections::HashSet<&str> =
605 self.view.column_order.iter().map(String::as_str).collect();
606 let joining: Vec<String> = dataset
607 .schema
608 .iter_names()
609 .map(|name| name.to_string())
610 .filter(|name| {
611 name != crate::formats::schema_union::DRIFT_COLUMN && !known.contains(name.as_str())
612 })
613 .collect();
614 drop(known);
615 self.view.column_order.extend(joining);
616 self.view
619 .column_order
620 .retain(|name| dataset.schema.contains(name.as_str()));
621 let schema = dataset.schema.clone();
622 match (remote, self.remote_files.as_mut()) {
625 (Some(found), Some(remote)) => {
626 remote.urls = Arc::new(found.urls);
627 remote.scan = found.scan;
628 remote.count = found.count;
631 }
632 (Some(found), None) => self.remote_files = Some(found.into()),
634 (None, _) => {}
635 }
636 self.record_dataset_schema(dataset, file_rows, files);
639 self.footers_pending = None;
640 self.replace_root(lf, schema);
643 self.pristine_rows = None;
645 self.view.observed_bytes_per_row = None;
648 if !row_groups.is_empty() {
651 self.record_file_row_groups(&row_groups);
652 }
653 self.deferred(Self::apply_transformations);
657 Ok(())
658 }
659
660 pub fn visible_lf(&self) -> LazyFrame {
664 Self::without_drift(self.view.lf.clone())
665 }
666
667 pub fn scans_a_temp_file(&self) -> bool {
670 self.decompress_temp_file.is_some() || !self.converted.is_empty()
671 }
672
673 pub fn scans_a_download(&self) -> bool {
676 self.download.is_some()
677 }
678
679 pub fn read_mode(&self) -> Option<crate::ReadMode> {
681 self.read_mode
682 }
683
684 pub fn read_as(&self) -> Option<crate::FileFormat> {
686 self.read_as
687 }
688
689 pub fn fetched(&self) -> bool {
691 self.fetched
692 }
693
694 pub(crate) fn temp_files(&self) -> Vec<&Path> {
697 let files = self.decompress_temp_file.iter().map(|file| file.path());
698 let files = files.chain(self.download.iter().map(|download| download.path()));
699 let files = files.chain(self.converted.iter().map(|file| file.path()));
700 files.collect()
701 }
702
703 pub(super) fn without_drift(lf: LazyFrame) -> LazyFrame {
705 lf.drop(by_name(
706 [crate::formats::schema_union::DRIFT_COLUMN],
707 false,
708 false,
709 ))
710 }
711
712 pub(super) fn query_source(&self) -> LazyFrame {
715 Self::without_drift(self.original_lf.clone())
716 }
717
718 pub(super) fn query_source_schema(&self) -> Arc<Schema> {
720 Self::without_source_rows(self.original_schema.clone()).0
721 }
722
723 pub fn drifts(&self) -> bool {
725 self.view.drift_column_present
726 }
727
728 pub fn drift_groups(&self) -> Arc<Vec<crate::formats::schema_union::DriftGroup>> {
730 self.view.drift_groups.clone()
731 }
732
733 pub fn can_name_source_files(&self) -> bool {
736 self.view.drift_column_present
737 && !self.drift_files.is_empty()
738 && self.drift_files.len() == self.drift_file_starts.len()
739 }
740
741 pub fn export_frame(&self, name_files: bool) -> ExportFrame {
745 if name_files && self.can_name_source_files() {
746 ExportFrame {
747 lf: self.view.lf.clone(),
748 files: Some(SourceFiles {
749 names: Arc::new(self.drift_files.clone()),
750 starts: Arc::new(self.drift_file_starts.clone()),
751 }),
752 }
753 } else {
754 ExportFrame {
755 lf: self.visible_lf(),
756 files: None,
757 }
758 }
759 }
760
761 pub fn notes(&self) -> Vec<crate::notes::Note> {
764 let mut notes = crate::notes::merged(
767 &self.open_notes,
768 &self.view.notes,
769 &self.view.view_notes,
770 self.dataset_schema.as_ref(),
771 );
772 if let Some(pushdown) = &self.pushdown {
774 notes.extend(pushdown.notes());
775 }
776 notes.extend(self.unfit_notes.iter().flatten().cloned());
777 notes.extend(self.view.changes_dropped.iter().cloned());
778 if let Some((version, unfit)) = &self.changes_unfit
779 && *version == self.view.changes_version
780 {
781 notes.extend(unfit.iter().cloned());
782 }
783 notes
784 }
785
786 pub(crate) fn window_for_quality(
789 &self,
790 scope: &crate::analysis::data_quality::QualityScope,
791 ) -> Option<Arc<dyn crate::formats::pushdown::Windowed>> {
792 use crate::analysis::data_quality::QualityScope;
793 matches!(scope, QualityScope::WholeSource | QualityScope::CurrentView)
794 .then(|| {
795 self.fixed_window
796 .clone()
797 .filter(|_| self.is_pristine() && self.indexing().is_none())
798 })
799 .flatten()
800 }
801
802 pub fn not_the_table(&self) -> Option<&'static str> {
805 self.not_the_table
806 }
807
808 pub fn other_tables(&self) -> &[String] {
810 &self.other_tables
811 }
812
813 pub fn format_read(&self) -> Option<&Arc<crate::formats::Read>> {
815 self.format_read.as_ref()
816 }
817
818 pub(super) fn window_now(&self) -> Option<Arc<dyn crate::formats::pushdown::Windowed>> {
821 if let Some(window) = self.follow_window() {
822 return Some(Arc::new(window));
823 }
824 if let Some(records) = self.fixed_window.as_ref().filter(|_| self.is_pristine()) {
825 return Some(records.clone());
826 }
827 self.pushed_view().map(|view| view.window)
828 }
829
830 pub(crate) fn pushed_view(&self) -> Option<crate::formats::pushdown::PushedView> {
833 let pushdown = self.pushdown.as_ref()?;
834 if !self.scan_is_the_root() || self.view.drift_column_present {
835 return None;
836 }
837 let sort: Vec<(String, bool)> = self
838 .view
839 .sort_columns
840 .iter()
841 .cloned()
842 .zip(self.view.sort_descending.iter().copied())
843 .collect();
844 pushdown.view(&self.view.filters, &sort, !self.view.sort_ascending)
845 }
846
847 pub(crate) fn source_counter(&self) -> Option<crate::formats::pushdown::Counter> {
850 if let Some(counter) = self.follow_counter() {
851 return Some(counter);
852 }
853 self.pushed_view().map(|view| view.counter)
854 }
855
856 fn follow_known(&self) -> Option<&[(usize, usize)]> {
859 self.follow_known
860 .as_ref()
861 .filter(|(generation, _)| *generation == self.view.len_generation)
862 .map(|(_, known)| known.as_slice())
863 }
864
865 fn follow_window(&self) -> Option<crate::loading::follow::Window> {
868 let follow = self.follow.as_ref()?;
869 let known = if self.is_pristine() {
870 None
871 } else if self.scan_is_the_root()
872 && self.view.sort_columns.is_empty()
873 && self.view.sort_ascending
874 {
875 Some(self.follow_known()?.to_vec())
876 } else {
877 return None;
878 };
879 Some(crate::loading::follow::Window {
880 lf: self.view.lf.clone(),
881 path: follow.path().to_path_buf(),
882 marks: follow.marks().clone(),
883 known,
884 })
885 }
886
887 fn follow_counter(&self) -> Option<crate::formats::pushdown::Counter> {
890 let follow = self.follow.as_ref()?;
891 let &(before, row) = self.follow_known()?.last()?;
892 let rest = crate::loading::follow::from_marks(
893 &self.view.lf,
894 follow.path(),
895 follow.marks(),
896 row,
897 None,
898 )?;
899 let streaming = self.polars_streaming;
900 Some(Arc::new(move || {
901 let df = crate::analysis::statistics::collect_lazy(row_count_lf(&rest), streaming)?;
902 let after = match df.get(0).and_then(|row| row.first().cloned()) {
903 Some(AnyValue::UInt64(n)) => n as usize,
904 _ => 0,
905 };
906 Ok(before + after)
907 }))
908 }
909
910 pub fn delimited_read(&self) -> Option<&Arc<crate::formats::delimited_spec::DelimitedRead>> {
913 self.delimited.as_ref()
914 }
915
916 pub fn unit_of(&self, column: &str) -> Option<&str> {
920 if self.delimited.is_none() && self.file_units.is_empty() {
921 return None;
922 }
923 let loaded = match &self.view.lineage {
924 None => column,
925 Some(lineage) => lineage
926 .iter()
927 .find(|(shown, _)| shown == column)
928 .map(|(_, loaded)| loaded.as_str())?,
929 };
930 match &self.delimited {
931 Some(read) => read.unit_of(loaded),
932 None => self
933 .file_units
934 .iter()
935 .find(|(name, _)| name == loaded)
936 .map(|(_, unit)| unit.as_str()),
937 }
938 }
939
940 pub fn units(&self) -> Vec<(String, String)> {
942 if self.delimited.is_none() && self.file_units.is_empty() {
943 return Vec::new();
944 }
945 self.view
946 .schema
947 .iter_names()
948 .filter_map(|name| Some((name.to_string(), self.unit_of(name)?.to_string())))
949 .collect()
950 }
951
952 pub fn format_detail(&self) -> Option<&crate::formats::text_formats::Detail> {
954 self.detail.as_deref()
955 }
956
957 pub(crate) fn ended_journal_to_describe(&mut self) -> Option<LazyFrame> {
960 let follow = self.follow.as_mut()?;
961 if follow.described
962 || follow.live()
963 || follow.behind()
964 || follow.spool().is_none()
965 || self.read_as != Some(crate::FileFormat::Journal)
966 {
967 return None;
968 }
969 follow.described = true;
970 Some(self.original_lf.clone())
971 }
972
973 pub(crate) fn set_format_detail(&mut self, detail: crate::formats::text_formats::Detail) {
974 self.detail = Some(Arc::new(detail));
975 }
976
977 pub fn has_notes(&self) -> bool {
979 !self.view.notes.is_empty()
982 || !self.open_notes.is_empty()
983 || self.unfit_notes.as_ref().is_some_and(|n| !n.is_empty())
984 || !self.view.changes_dropped.is_empty()
985 || self
986 .changes_unfit
987 .as_ref()
988 .is_some_and(|(v, n)| *v == self.view.changes_version && !n.is_empty())
989 || self
990 .pushdown
991 .as_ref()
992 .is_some_and(|p| !p.notes().is_empty())
993 }
994
995 pub fn notes_unseen(&self) -> bool {
997 self.has_notes() && !self.view.notes_seen
998 }
999
1000 fn unread_row_runs(&self, column: &str) -> Vec<(usize, usize)> {
1003 let Some(dataset) = self.dataset_schema.as_ref() else {
1004 return Vec::new();
1005 };
1006 let conflicts: Vec<bool> = (0..self.drift_file_starts.len())
1007 .map(|file| {
1008 dataset
1009 .file_group
1010 .get(file)
1011 .and_then(|group| dataset.groups.get(*group as usize))
1012 .is_some_and(|group| group.unread.iter().any(|name| name == column))
1013 })
1014 .collect();
1015 conflicting_row_runs(&self.drift_file_starts, self.drift_dataset_rows, &conflicts)
1016 }
1017
1018 fn view_columns_with_conflicts(&self) -> Vec<crate::formats::schema_union::ColumnDrift> {
1021 let Some(dataset) = self.dataset_schema.as_ref() else {
1022 return Vec::new();
1023 };
1024 let named: HashSet<&str> = self
1025 .view
1026 .filters
1027 .iter()
1028 .map(|filter| filter.column.as_str())
1029 .chain(self.view.sort_columns.iter().map(String::as_str))
1030 .collect();
1031 dataset
1032 .columns
1033 .iter()
1034 .filter(|column| column.conflicting_files > 0 && named.contains(column.name.as_str()))
1035 .cloned()
1036 .collect()
1037 }
1038
1039 pub(super) fn view_exclusions(&self) -> Vec<(Vec<(usize, usize)>, crate::notes::Note)> {
1043 if !self.view.drift_column_present {
1044 return Vec::new();
1045 }
1046 let Some(dataset) = self.dataset_schema.as_ref() else {
1047 return Vec::new();
1048 };
1049 let mut out = Vec::new();
1050 for column in self.view_columns_with_conflicts() {
1051 let runs = self.unread_row_runs(&column.name);
1052 let rows: usize = runs.iter().map(|(start, end)| end - start).sum();
1053 if rows == 0 {
1054 continue;
1055 }
1056 let filtered = self
1057 .view
1058 .filters
1059 .iter()
1060 .any(|filter| filter.column.as_str() == column.name.as_str());
1061 let sorted = self
1062 .view
1063 .sort_columns
1064 .iter()
1065 .any(|sorted| sorted.as_str() == column.name.as_str());
1066 out.push((
1067 runs,
1068 crate::notes::left_out_note(&column, dataset, rows, filtered, sorted),
1069 ));
1070 }
1071 out
1072 }
1073
1074 pub(super) fn view_notes_only(&self) -> Vec<crate::notes::Note> {
1077 self.view_exclusions()
1078 .into_iter()
1079 .map(|(_, note)| note)
1080 .collect()
1081 }
1082
1083 pub(super) fn leave_out_unread_rows(
1084 &self,
1085 mut lf: LazyFrame,
1086 ) -> (LazyFrame, Vec<crate::notes::Note>) {
1087 let mut notes = Vec::new();
1088 for (runs, note) in self.view_exclusions() {
1089 let keep = runs
1090 .iter()
1091 .map(|(start, end)| {
1092 col(crate::formats::schema_union::DRIFT_COLUMN)
1093 .lt(lit(*start as u32))
1094 .or(col(crate::formats::schema_union::DRIFT_COLUMN).gt_eq(lit(*end as u32)))
1095 })
1096 .reduce(Expr::and);
1097 if let Some(keep) = keep {
1098 lf = lf.filter(keep);
1099 }
1100 notes.push(note);
1101 }
1102 (lf, notes)
1103 }
1104
1105 pub fn mark_notes_seen(&mut self) {
1107 self.view.notes_seen = true;
1108 }
1109
1110 pub fn read_column_as_text(&mut self, column: &str) -> PolarsResult<bool> {
1116 let name = PlSmallStr::from(column);
1117 let Some(dataset) = self.dataset_at_open.clone() else {
1118 return Ok(false);
1119 };
1120 if !self.view.drift_column_present || self.read_as_text.contains(&name) {
1121 return Ok(false);
1122 }
1123 if !dataset
1124 .columns
1125 .iter()
1126 .any(|drift| drift.name == name && drift.can_read_as_text())
1127 {
1128 return Ok(false);
1129 }
1130
1131 let mut as_text = self.read_as_text.clone();
1132 as_text.push(name);
1133
1134 let drift = crate::formats::schema_union::ScanDrift::new(
1137 &self.drift_files,
1138 &dataset,
1139 &self.file_rows(),
1140 );
1141 let scanned = match self.remote_files.as_ref() {
1142 Some(remote) => (remote.scan)(&remote.urls, &as_text),
1143 None => crate::formats::schema_union::lenient_scan(
1144 &self.drift_files,
1145 dataset.schema.clone(),
1146 None,
1147 drift.as_ref(),
1148 &as_text,
1149 ),
1150 };
1151 let lf = match scanned {
1152 Ok(lf) => lf,
1153 Err(e) => {
1154 self.error = Some(e.clone());
1156 return Err(e);
1157 }
1158 };
1159 let lf = if self.remote_files.is_some() {
1162 lf
1163 } else {
1164 crate::loading::open_scan::hoist_partition_columns(
1165 lf,
1166 &dataset.schema,
1167 self.partition_columns.as_deref().unwrap_or(&[]),
1168 drift.is_some(),
1169 )
1170 };
1171
1172 let view = dataset.reading_as_text(&as_text);
1173 self.read_as_text = as_text;
1174 let schema = view.schema.clone();
1176 self.view.drift_groups = Arc::new(view.groups.clone());
1177 self.groups_at_open = self.view.drift_groups.clone();
1178 self.view.notes = Self::notes_datui_can_act_on(&view, self.view.drift_column_present);
1179 self.notes_at_open = self.view.notes.clone();
1180 self.view.notes_seen = false;
1182 self.dataset_schema = Some(view);
1183 self.replace_root(lf, schema);
1184 self.apply_transformations();
1186 Ok(true)
1187 }
1188
1189 fn notes_datui_can_act_on(
1193 dataset: &crate::formats::schema_union::DatasetSchema,
1194 counted: bool,
1195 ) -> Vec<crate::notes::Note> {
1196 let mut notes = crate::notes::from_dataset(dataset);
1197 if !counted {
1198 for note in &mut notes {
1199 note.read_as_text = None;
1200 }
1201 }
1202 notes
1203 }
1204
1205 pub(super) fn file_rows(&self) -> Vec<usize> {
1208 self.drift_file_starts
1209 .iter()
1210 .enumerate()
1211 .map(|(file, start)| {
1212 self.drift_file_starts
1213 .get(file + 1)
1214 .copied()
1215 .unwrap_or(self.drift_dataset_rows)
1216 .saturating_sub(*start)
1217 })
1218 .collect()
1219 }
1220
1221 pub fn read_as_text(&self) -> &[PlSmallStr] {
1223 &self.read_as_text
1224 }
1225
1226 pub fn dataset_schema(&self) -> Option<&crate::formats::schema_union::DatasetSchema> {
1228 self.dataset_schema.as_ref()
1229 }
1230
1231 pub fn remote_files_counter(&self) -> Option<FileCounter> {
1235 if self.footers_pending.is_some() {
1236 return None;
1237 }
1238 self.remote_files
1239 .as_ref()
1240 .filter(|f| f.offsets.is_none() && self.is_pristine())
1241 .map(|f| f.count.clone())
1242 }
1243
1244 pub(super) fn record_file_row_groups(&mut self, groups: &[Vec<usize>]) {
1247 if let Some(files) = self.remote_files.as_mut() {
1248 if groups.len() != files.urls.len() {
1249 return;
1250 }
1251 let mut offsets = Vec::with_capacity(groups.len() + 1);
1252 offsets.push(0);
1253 for file in groups {
1254 offsets.push(offsets.last().unwrap_or(&0) + file.iter().sum::<usize>());
1255 }
1256 files.offsets = Some(offsets);
1257 }
1258 let flat: Vec<usize> = groups.iter().flatten().copied().collect();
1259 if self.is_pristine() {
1260 self.record_row_groups(&flat);
1261 } else {
1262 let mut row_offsets = Vec::with_capacity(flat.len() + 1);
1264 row_offsets.push(0);
1265 for n in &flat {
1266 row_offsets.push(row_offsets.last().unwrap_or(&0) + n);
1267 }
1268 self.row_group_offsets = Some(row_offsets);
1269 }
1270 }
1271}