1use super::*;
5
6pub struct CollectRequest {
8 pub lf: LazyFrame,
10 pub polars_streaming: bool,
12 pub buffer_start: usize,
14 pub buffer_end: usize,
16 pub plan: FillPlan,
18}
19
20pub struct FillPlan {
25 buffer_start: usize,
26 buffer_end: usize,
27 num_rows: usize,
28 count_known: bool,
29 indexing: bool,
32 held: Option<(DataFrame, usize)>,
35 view_start: usize,
36 view_len: usize,
37 max_rows: usize,
38 max_mb: usize,
39}
40
41impl FillPlan {
42 pub fn fit(mut self, df: DataFrame) -> CollectResult {
45 let returned = df.height();
46 let bytes_per_row = (returned > 0).then(|| (df.estimated_size() / returned).max(1));
47 let (df, start, seam) = match self.held.take() {
49 Some((mut held, held_start)) if held_start + held.height() == self.buffer_start => {
50 let seam = held.height();
51 match held.vstack_mut(&df) {
52 Ok(_) => (held, held_start, Some(seam)),
53 Err(_) => (df, self.buffer_start, None),
54 }
55 }
56 Some((held, held_start))
57 if returned > 0 && self.buffer_start + returned == held_start =>
58 {
59 match df.vstack(&held) {
60 Ok(joined) => (joined, self.buffer_start, Some(returned)),
61 Err(_) => (df, self.buffer_start, None),
62 }
63 }
64 _ => (df, self.buffer_start, None),
65 };
66 let (df, start) = self.cut_to_caps(df, start, seam);
67 let df = fewer_chunks(df);
68 CollectResult {
69 df,
70 start,
71 returned,
72 bytes_per_row,
73 buffer_start: self.buffer_start,
74 buffer_end: self.buffer_end,
75 num_rows: self.num_rows,
76 count_known: self.count_known,
77 indexing: self.indexing,
78 }
79 }
80
81 fn cut_to_caps(&self, df: DataFrame, start: usize, seam: Option<usize>) -> (DataFrame, usize) {
86 let total = df.height();
87 if total == 0 {
88 return (df, start);
89 }
90 let mut max_rows = total;
92 if self.max_rows > 0 {
93 max_rows = max_rows.min(self.max_rows);
94 }
95 if self.max_mb > 0 {
96 let bytes_per_row = (df.estimated_size() / total).max(1);
97 max_rows = max_rows.min(self.max_mb * 1024 * 1024 / bytes_per_row);
98 }
99 let max_rows = max_rows.max(1);
100 if max_rows >= total {
101 return (df, start);
102 }
103 let view_off = self.view_start.saturating_sub(start).min(total);
104 let view_len = self.view_len.max(1).min(total);
105 let view_center = view_off + view_len / 2;
106 let mut keep_start = view_center.saturating_sub(max_rows / 2);
107 if keep_start + max_rows > total {
108 keep_start = total - max_rows;
109 }
110 let kept = max_rows.min(total - keep_start);
111 (trim_rows(df, keep_start, kept, seam), start + keep_start)
112 }
113}
114
115pub struct CollectResult {
118 pub(super) df: DataFrame,
120 pub(super) start: usize,
122 returned: usize,
124 bytes_per_row: Option<usize>,
126 buffer_start: usize,
128 buffer_end: usize,
129 num_rows: usize,
131 count_known: bool,
134 indexing: bool,
136}
137
138impl CollectResult {
139 pub(crate) fn rows(&self) -> &DataFrame {
141 &self.df
142 }
143}
144
145pub(super) const STRING_BYTES_GUESS: usize = 40;
147
148pub(super) fn estimate_bytes_per_row(
152 schema: &Schema,
153 columns: &[String],
154 column_bytes: &[(String, usize)],
155) -> usize {
156 let footer_width = |name: &String| {
157 column_bytes
158 .iter()
159 .find(|(n, _)| n == name)
160 .map(|(_, w)| *w)
161 };
162 columns
163 .iter()
164 .map(|name| match schema.get(name.as_str()) {
165 Some(DataType::String) => 16 + footer_width(name).unwrap_or(STRING_BYTES_GUESS - 16),
166 Some(DataType::Binary) => 16 + binary_stub().len(),
167 Some(DataType::Boolean) => 1,
168 Some(DataType::Null) => 0,
169 Some(dtype) if dtype.is_primitive_numeric() || dtype.is_temporal() => {
170 match dtype.to_physical() {
171 DataType::Int8 | DataType::UInt8 => 1,
172 DataType::Int16 | DataType::UInt16 => 2,
173 DataType::Int32 | DataType::UInt32 | DataType::Float32 => 4,
174 DataType::Int128 => 16,
175 _ => 8,
176 }
177 }
178 Some(DataType::Decimal(..)) => 16,
179 _ => footer_width(name).unwrap_or(64),
180 })
181 .sum::<usize>()
182 .max(1)
183}
184
185pub(super) const MAX_BUFFER_CHUNKS: usize = 16;
188
189fn fewer_chunks(df: DataFrame) -> DataFrame {
192 let chunks = df
193 .columns()
194 .iter()
195 .filter_map(Column::as_series)
196 .map(|s| s.chunks().len())
197 .max()
198 .unwrap_or(0);
199 if chunks <= MAX_BUFFER_CHUNKS {
200 return df;
201 }
202 let rows = df.height();
203 compact_rows(df, 0, rows, None)
204}
205
206pub(super) fn trim_rows(
211 df: DataFrame,
212 offset: usize,
213 len: usize,
214 seam: Option<usize>,
215) -> DataFrame {
216 if backing_rows(&df, offset, len) > len + len / 4 {
217 compact_rows(df, offset, len, seam)
218 } else {
219 df.slice(offset as i64, len)
220 }
221}
222
223pub(super) fn backing_rows(df: &DataFrame, offset: usize, len: usize) -> usize {
226 let end = offset + len;
227 df.columns()
228 .iter()
229 .filter_map(Column::as_series)
230 .map(|s| {
231 let mut start = 0;
232 let mut touched = 0;
233 for chunk in s.chunks() {
234 let chunk_end = start + chunk.len();
235 if start < end && offset < chunk_end {
236 touched += chunk.len();
237 }
238 start = chunk_end;
239 }
240 touched
241 })
242 .max()
243 .unwrap_or(len)
244}
245
246pub(super) fn compact_rows(
252 df: DataFrame,
253 offset: usize,
254 len: usize,
255 seam: Option<usize>,
256) -> DataFrame {
257 use polars::series::builder::SeriesBuilder;
258 use polars_arrow::array::builder::ShareStrategy;
259 #[cfg(test)]
260 tests::COMPACTIONS.with(|count| count.set(count.get() + 1));
261 let len = len.min(df.height().saturating_sub(offset));
262 let pieces = match seam.filter(|&seam| offset < seam && seam < offset + len) {
263 Some(seam) => vec![(offset, seam - offset), (seam, offset + len - seam)],
264 None => vec![(offset, len)],
265 };
266 let copy = |series: &Series, (offset, len): (usize, usize)| {
267 let mut builder = SeriesBuilder::new(series.dtype().clone());
268 builder.reserve(len);
269 builder.subslice_extend(series, offset, len, ShareStrategy::Never);
270 builder.freeze(series.name().clone())
271 };
272 let columns = df
273 .into_columns()
274 .into_iter()
275 .map(|column| match column {
276 Column::Scalar(constant) => {
277 Column::new_scalar(constant.name().clone(), constant.scalar().clone(), len)
278 }
279 Column::Series(series) => {
280 let mut kept = copy(&series, pieces[0]);
281 for &piece in &pieces[1..] {
282 if kept.append_owned(copy(&series, piece)).is_err() {
283 kept = copy(&series, (offset, len));
284 break;
285 }
286 }
287 kept.into_column()
288 }
289 })
290 .collect();
291 DataFrame::new(len, columns).unwrap_or_else(|_| DataFrame::empty_with_height(len))
293}
294
295pub(super) fn shrink_around_view(
298 view_start: usize,
299 view_end: usize,
300 max_len: usize,
301 floor: usize,
302 ceil: usize,
303 buffer_start: &mut usize,
304 buffer_end: &mut usize,
305) {
306 if buffer_end.saturating_sub(*buffer_start) <= max_len {
307 return;
308 }
309 let view_len = view_end.saturating_sub(view_start);
310 if view_len >= max_len {
311 *buffer_start = view_start;
312 *buffer_end = (view_start + max_len).min(ceil);
313 return;
314 }
315 let half = (max_len - view_len) / 2;
316 *buffer_end = (view_end + half).min(ceil);
317 *buffer_start = buffer_end.saturating_sub(max_len).max(floor);
318 if *buffer_start > view_start {
319 *buffer_start = view_start;
320 }
321 *buffer_end = (*buffer_start + max_len).min(ceil);
322}
323
324pub(super) const MAX_FILES_PER_BUFFER: usize = 16;
326
327pub(super) fn limit_files(
330 offsets: &[usize],
331 view_start: usize,
332 view_end: usize,
333 start: usize,
334 end: usize,
335 max_files: usize,
336) -> (usize, usize) {
337 let (Some((first, last)), Some((view_first, view_last))) = (
338 files_holding(offsets, start, end.saturating_sub(start)),
339 files_holding(
340 offsets,
341 view_start,
342 view_end.saturating_sub(view_start).max(1),
343 ),
344 ) else {
345 return (start, end);
346 };
347 let opened = |from: usize, to: usize| (from..=to).filter(|&i| holds_rows(offsets, i)).count();
349 if opened(first, last) <= max_files {
350 return (start, end);
351 }
352 let (mut lo, mut hi) = (view_first.max(first), view_last.min(last));
353 let mut files = opened(lo, hi);
354 while files < max_files && (hi < last || lo > first) {
355 if hi < last {
356 hi += 1;
357 files += usize::from(holds_rows(offsets, hi));
358 }
359 if files < max_files && lo > first {
360 lo -= 1;
361 files += usize::from(holds_rows(offsets, lo));
362 }
363 }
364 (start.max(offsets[lo]), end.min(offsets[hi + 1]))
365}
366
367pub(super) fn holds_rows(offsets: &[usize], i: usize) -> bool {
369 offsets[i + 1] > offsets[i]
370}
371
372pub(super) fn files_with_rows(offsets: &[usize], first: usize, last: usize) -> Vec<usize> {
375 (first..=last).filter(|&i| holds_rows(offsets, i)).collect()
376}
377
378pub(super) fn files_holding(offsets: &[usize], start: usize, len: usize) -> Option<(usize, usize)> {
381 let files = offsets.len().checked_sub(1)?;
382 let total = *offsets.last()?;
383 if files == 0 || len == 0 || start >= total {
384 return None;
385 }
386 let end = (start + len).min(total);
387 let file_of = |row: usize| offsets.partition_point(|&o| o <= row).saturating_sub(1);
390 Some((file_of(start), file_of(end - 1).min(files - 1)))
391}
392
393#[derive(Clone, Copy, Default)]
395pub(super) struct Sources<'a> {
396 pub(super) files: Option<&'a RemoteFiles>,
398 pub(super) records: Option<&'a dyn crate::formats::pushdown::Windowed>,
400 pub(super) csv: Option<&'a Arc<super::csv_marks::CsvMarks>>,
402}
403
404pub(super) fn window_of(
407 lf: &LazyFrame,
408 sources: Sources<'_>,
409 read_as_text: &[PlSmallStr],
410 start: usize,
411 len: usize,
412 all_columns: Vec<Expr>,
413) -> PolarsResult<LazyFrame> {
414 let Sources {
415 files,
416 records,
417 csv,
418 } = sources;
419 if let Some(records) = records {
422 return Ok(records.window(start, len)?.select(all_columns));
423 }
424 if let Some(window) = csv.and_then(|marks| marks.window(lf, start, len)) {
425 return Ok(window.select(all_columns));
426 }
427 if let Some((files, offsets)) = files.and_then(|f| f.offsets.as_ref().map(|o| (f, o)))
428 && let Some((first, last)) = files_holding(offsets, start, len)
429 {
430 let urls: Vec<String> = files_with_rows(offsets, first, last)
433 .into_iter()
434 .map(|i| files.urls[i].clone())
435 .collect();
436 let lf = (files.scan)(&urls, read_as_text)?;
437 return Ok(lf
438 .select(all_columns)
439 .slice((start - offsets[first]) as i64, len as u32));
440 }
441 Ok(lf
442 .clone()
443 .select(all_columns)
444 .slice(start as i64, len as u32))
445}
446
447#[derive(Clone)]
450pub(crate) struct ViewRows {
451 lf: LazyFrame,
452 files: Option<RemoteFiles>,
453 records: Option<Arc<dyn crate::formats::pushdown::Windowed>>,
455 csv: Option<Arc<super::csv_marks::CsvMarks>>,
457 read_as_text: Vec<PlSmallStr>,
458 pub(crate) buffer: Option<(DataFrame, usize)>,
460 pub(crate) num_rows: Option<usize>,
462 pub(crate) streaming: bool,
463 pub(crate) whole: bool,
465 pub(crate) reads_up_to: bool,
467}
468
469pub(crate) fn sees_every_row_first(lf: &LazyFrame) -> bool {
472 use polars::lazy::dsl::DslPlan;
473 lf.logical_plan.into_iter().any(|node| {
474 matches!(
475 node,
476 DslPlan::Sort { .. } | DslPlan::GroupBy { .. } | DslPlan::Pivot { .. }
477 )
478 })
479}
480
481pub(crate) fn reads_up_to_a_window(lf: &LazyFrame) -> bool {
484 use polars::lazy::dsl::{DslPlan, FileScanDsl};
485 lf.logical_plan.into_iter().any(|node| match node {
486 DslPlan::Filter { .. } => true,
487 DslPlan::Scan { scan_type, .. } => !matches!(
488 **scan_type,
489 FileScanDsl::Parquet { .. } | FileScanDsl::Ipc { .. }
490 ),
491 _ => false,
492 })
493}
494
495impl ViewRows {
496 pub(crate) fn window(
498 &self,
499 start: usize,
500 len: usize,
501 exprs: Vec<Expr>,
502 ) -> PolarsResult<LazyFrame> {
503 window_of(
504 &self.lf,
505 Sources {
506 files: self.files.as_ref(),
507 records: self.records.as_deref(),
508 csv: self.csv.as_ref(),
509 },
510 &self.read_as_text,
511 start,
512 len,
513 exprs,
514 )
515 }
516
517 #[cfg(test)]
519 pub(crate) fn of(lf: LazyFrame, buffer: Option<(DataFrame, usize)>) -> Self {
520 Self {
521 whole: sees_every_row_first(&lf),
522 reads_up_to: reads_up_to_a_window(&lf),
523 lf,
524 files: None,
525 records: None,
526 csv: None,
527 read_as_text: Vec::new(),
528 buffer,
529 num_rows: None,
530 streaming: false,
531 }
532 }
533}
534
535pub(super) const PAGES_DECODED_AHEAD: usize = 8192;
538
539pub(super) const BYTES_DECODED_AHEAD: usize = 8 << 20;
542
543pub(super) fn decodes_pages(lf: &LazyFrame) -> bool {
545 use polars::lazy::dsl::{DslPlan, FileScanDsl, ScanSources};
546 let remote = |path: &str| crate::cloud::source::is_remote_url(Path::new(path));
547 (&lf.logical_plan).into_iter().any(|node| match node {
548 DslPlan::Scan {
549 sources: ScanSources::Paths(paths),
550 scan_type,
551 ..
552 } => {
553 matches!(**scan_type, FileScanDsl::Parquet { .. })
554 && !paths.iter().any(|path| remote(path.as_str()))
555 }
556 _ => false,
557 })
558}
559
560pub(super) fn align_to_row_groups(
565 offsets: &[usize],
566 view_start: usize,
567 view_end: usize,
568 start: usize,
569 end: usize,
570 cap: usize,
571) -> (usize, usize) {
572 let Some(groups) = offsets.len().checked_sub(1).filter(|n| *n > 0) else {
573 return (start, end);
574 };
575 let group_of = |row: usize| {
576 offsets
577 .partition_point(|&o| o <= row)
578 .saturating_sub(1)
579 .min(groups - 1)
580 };
581 let last_row = |s: usize, e: usize| e.saturating_sub(1).max(s);
582 let (mut lo, mut hi) = (
583 group_of(view_start),
584 group_of(last_row(view_start, view_end)),
585 );
586 let (want_lo, want_hi) = (group_of(start), group_of(last_row(start, end)));
587 let fits = |lo: usize, hi: usize| cap == 0 || offsets[hi + 1] - offsets[lo] <= cap;
588 loop {
589 if hi < want_hi && fits(lo, hi + 1) {
590 hi += 1;
591 } else if lo > want_lo && fits(lo - 1, hi) {
592 lo -= 1;
593 } else {
594 break;
595 }
596 }
597 (offsets[lo], offsets[hi + 1])
598}
599
600impl DataTableState {
601 pub fn scroll_would_trigger_collect(&self, rows: i64) -> bool {
604 if rows < 0 && self.view.start_row == 0 {
605 return false;
606 }
607 let new_start_row = if self.view.start_row as i64 + rows <= 0 {
608 0
609 } else {
610 if let Some(df) = self.view.df.as_ref()
611 && rows > 0
612 && df.shape().0 <= self.visible_rows
613 {
614 return false;
615 }
616 let unclamped = (self.view.start_row as i64 + rows) as usize;
617 if rows > 0 {
618 unclamped.min(self.view.num_rows.saturating_sub(self.visible_rows))
619 } else {
620 unclamped
621 }
622 };
623 if new_start_row == self.view.start_row {
624 return false;
625 }
626 let view_end = new_start_row
627 + self
628 .visible_rows
629 .min(self.view.num_rows.saturating_sub(new_start_row));
630 let within_buffer = new_start_row >= self.view.buffered_start_row
631 && view_end <= self.view.buffered_end_row
632 && self.view.buffered_end_row > 0;
633 !within_buffer
634 }
635
636 pub fn slide_table(&mut self, rows: i64) -> bool {
639 if rows < 0 && self.view.start_row == 0 {
640 return false;
641 }
642
643 let new_start_row = if self.view.start_row as i64 + rows <= 0 {
644 0
645 } else {
646 if let Some(df) = self.view.df.as_ref()
647 && rows > 0
648 && df.shape().0 <= self.visible_rows
649 {
650 return false;
651 }
652 let unclamped = (self.view.start_row as i64 + rows) as usize;
653 if rows > 0 {
654 unclamped.min(self.view.num_rows.saturating_sub(self.visible_rows))
657 } else {
658 unclamped
659 }
660 };
661
662 if new_start_row == self.view.start_row {
663 return false;
664 }
665
666 let view_end = new_start_row
667 + self
668 .visible_rows
669 .min(self.view.num_rows.saturating_sub(new_start_row));
670 let within_buffer = new_start_row >= self.view.buffered_start_row
671 && view_end <= self.view.buffered_end_row
672 && self.view.buffered_end_row > 0;
673
674 self.view.start_row = new_start_row;
675
676 if within_buffer {
677 if self.table_state.selected().is_none() {
678 self.table_state.select(Some(0));
679 }
680 false
681 } else {
682 true }
684 }
685
686 #[cfg(test)]
688 pub fn collect(&mut self) {
689 if self.defer_collect {
690 return;
691 }
692 if !self.view.num_rows_valid {
693 match collect_lazy(row_count_lf(&self.view.lf), self.polars_streaming) {
696 Ok(df) => {
697 self.error = None;
698 let n = match df.get(0).as_deref().and_then(|row| row.first()) {
699 Some(AnyValue::UInt64(len)) => *len as usize,
700 _ => 0,
701 };
702 self.set_num_rows(n);
703 }
704 Err(e) => {
705 self.error = Some(e);
706 self.set_num_rows(0);
707 }
708 }
709 }
710 let Some(request) = self.prepare_async_collect(None) else {
711 return;
712 };
713 match collect_lazy(request.lf, request.polars_streaming) {
714 Ok(df) => self.apply_async_collect(request.plan.fit(df)),
715 Err(e) => self.error = Some(e),
716 }
717 }
718
719 #[cfg(not(test))]
722 pub(super) fn collect(&mut self) {
723 if !self.defer_collect {
724 self.needs_recollect = true;
725 }
726 }
727
728 pub(crate) fn binary_stub_exprs(&self) -> Vec<Expr> {
732 self.view
733 .column_order
734 .iter()
735 .map(|name| {
736 if matches!(self.view.schema.get(name.as_str()), Some(DataType::Binary)) {
737 lit(binary_stub()).alias(name.as_str())
738 } else {
739 col(name.as_str())
740 }
741 })
742 .collect()
743 }
744
745 pub fn prepare_async_collect(
749 &mut self,
750 num_rows_override: Option<usize>,
751 ) -> Option<CollectRequest> {
752 if self.visible_rows > 0 {
753 self.proximity_threshold = self.proximity();
754 }
755
756 if let Some(n) = num_rows_override {
757 self.view.num_rows = n;
758 self.view.num_rows_valid = true;
759 }
760
761 let count_known = self.view.num_rows_valid;
764 let bound = self.num_rows_bound();
765
766 if count_known {
767 if self.view.num_rows > 0 {
768 let max_start = self.view.num_rows.saturating_sub(1);
769 if self.view.start_row > max_start {
770 self.view.start_row = max_start;
771 }
772 } else {
773 self.view.start_row = 0;
775 self.drop_buffer();
776 self.view.df = None;
777 self.view.locked_df = None;
778 return None;
779 }
780 }
781
782 if self.view.column_order.is_empty() {
785 self.drop_buffer();
786 self.view.df = None;
787 self.view.locked_df = None;
788 return None;
789 }
790
791 let view_start = self.view.start_row;
792 let view_end = self.view.start_row + self.visible_rows.min(bound - self.view.start_row);
793 let within_buffer = view_start >= self.view.buffered_start_row
794 && view_end <= self.view.buffered_end_row
795 && self.view.buffered_end_row > 0;
796
797 let (new_buffer_start, new_buffer_end) = if within_buffer {
799 let dist_to_start = view_start.saturating_sub(self.view.buffered_start_row);
800 let dist_to_end = self.view.buffered_end_row.saturating_sub(view_end);
801 let needs_expansion_back =
802 dist_to_start <= self.proximity_behind() && self.view.buffered_start_row > 0;
803 let needs_expansion_forward =
804 dist_to_end <= self.proximity_threshold && self.view.buffered_end_row < bound;
805
806 if !needs_expansion_back && !needs_expansion_forward {
807 (self.view.buffered_start_row, self.view.buffered_end_row)
809 } else {
810 let mut s = if needs_expansion_back {
811 view_start.saturating_sub(self.reach_behind())
812 } else {
813 self.view.buffered_start_row
814 };
815 let mut e = if needs_expansion_forward {
816 (view_end + self.reach_ahead()).min(bound)
817 } else {
818 self.view.buffered_end_row
819 };
820 self.fit_window(view_start, view_end, &mut s, &mut e);
821 (s, e)
822 }
823 } else {
824 let had_buffer = self.view.buffered_end_row > 0;
825 let scrolled_past_end = had_buffer && view_start >= self.view.buffered_end_row;
826 let scrolled_past_start = had_buffer && view_end <= self.view.buffered_start_row;
827 let extend_forward_ok = scrolled_past_end
828 && (view_start - self.view.buffered_end_row) <= self.reach_ahead();
829 let extend_backward_ok = scrolled_past_start
830 && (self.view.buffered_start_row - view_end) <= self.reach_behind();
831
832 let mut s;
833 let mut e;
834 if extend_forward_ok {
835 s = self.view.buffered_start_row;
836 e = (view_end + self.reach_ahead()).min(bound);
837 } else if extend_backward_ok {
838 s = view_start.saturating_sub(self.reach_behind());
839 e = self.view.buffered_end_row;
840 } else {
841 s = view_start.saturating_sub(self.reach_behind());
842 e = (view_end + self.reach_ahead()).min(bound);
843 let min_initial_len = self.min_buffer_len();
844 let current_len = e.saturating_sub(s);
845 if current_len < min_initial_len {
846 let need = min_initial_len.saturating_sub(current_len);
847 let can_extend_end = bound.saturating_sub(e);
848 let can_extend_start = s;
849 if can_extend_end >= need {
850 e = (e + need).min(bound);
851 } else if can_extend_start >= need {
852 s = s.saturating_sub(need);
853 } else {
854 e = (e + can_extend_end).min(bound);
855 s = s.saturating_sub(need.saturating_sub(can_extend_end));
856 }
857 }
858 }
859 self.fit_window(view_start, view_end, &mut s, &mut e);
860 (s, e)
861 };
862
863 let buffer_size = new_buffer_end.saturating_sub(new_buffer_start);
864 if buffer_size == 0 {
865 return None;
866 }
867 if self.holds_buffer(new_buffer_start, new_buffer_end) {
870 self.slice_buffer_into_display();
871 if self.table_state.selected().is_none() {
872 self.table_state.select(Some(0));
873 }
874 return None;
875 }
876
877 let lf = match self.buffer_lf(new_buffer_start, buffer_size) {
878 Ok(lf) => lf,
879 Err(e) => {
880 self.error = Some(e);
881 return None;
882 }
883 };
884
885 let num_rows = if count_known {
888 self.view.num_rows
889 } else {
890 new_buffer_end
891 };
892 Some(CollectRequest {
893 lf,
894 polars_streaming: self.polars_streaming,
895 buffer_start: new_buffer_start,
896 buffer_end: new_buffer_end,
897 plan: self.fill_plan(new_buffer_start, new_buffer_end, num_rows, count_known),
898 })
899 }
900
901 pub(super) fn fill_plan(
904 &self,
905 buffer_start: usize,
906 buffer_end: usize,
907 num_rows: usize,
908 count_known: bool,
909 ) -> FillPlan {
910 let held = self
911 .abuts_buffer(buffer_start, buffer_end.saturating_sub(buffer_start))
912 .then(|| self.view.buffered_df.clone())
913 .flatten()
914 .map(|df| (df, self.view.buffered_start_row));
915 FillPlan {
916 buffer_start,
917 buffer_end,
918 num_rows,
919 count_known,
920 indexing: self.indexing().is_some(),
921 held,
922 view_start: self.view.start_row,
923 view_len: self.visible_rows,
924 max_rows: self.max_buffered_rows,
925 max_mb: self.max_buffered_mb,
926 }
927 }
928
929 pub fn apply_async_collect(&mut self, result: CollectResult) {
932 let CollectResult {
933 df,
934 start,
935 returned: returned_rows,
936 bytes_per_row,
937 buffer_start,
938 buffer_end,
939 num_rows,
940 count_known,
941 indexing,
942 } = result;
943 let requested_rows = buffer_end.saturating_sub(buffer_start);
944
945 if count_known {
946 self.view.num_rows = num_rows;
947 self.view.num_rows_valid = true;
948 } else if returned_rows < requested_rows
949 && (buffer_start == 0 || returned_rows > 0)
950 && !indexing
952 && self.indexing().is_none()
953 {
954 self.view.num_rows = buffer_start + returned_rows;
957 self.view.num_rows_valid = true;
958 } else if !self.view.num_rows_valid {
959 self.view.num_rows = self.view.num_rows.max(buffer_end);
962 }
963 self.error = None;
966 self.remember_pristine_count();
967
968 if bytes_per_row.is_some() {
969 self.view.observed_bytes_per_row = bytes_per_row;
970 }
971 let end = start + df.height();
975 let view_end = self.view.start_row + self.visible_rows.max(1);
976 let reaches_end = end >= buffer_start + returned_rows;
977 let shows_view = start <= self.view.start_row
978 && (self.view.start_row < end || (returned_rows < requested_rows && reaches_end));
979 if !shows_view {
980 self.needs_recollect = true;
981 return;
982 }
983 self.release_display_buffer();
984 self.view.buffered_start_row = start;
985 self.view.buffered_end_row = end;
986 self.view.buffered_df = Some(df);
987 self.slice_buffer_into_display();
989 if self.table_state.selected().is_none() {
990 self.table_state.select(Some(0));
991 }
992 if view_end > end && end < self.view.num_rows {
993 self.needs_recollect = true;
994 }
995 }
996
997 fn abuts_buffer(&self, start: usize, rows: usize) -> bool {
1000 self.stitches_buffer()
1001 && (start == self.view.buffered_end_row || start + rows == self.view.buffered_start_row)
1002 }
1003
1004 pub(crate) fn sampled_from(
1008 source: DataTableState,
1009 sample: crate::analysis::sampling::Sample,
1010 schema: &Schema,
1011 rows: Arc<crate::analysis::table_sample::SampleRows>,
1012 through: bool,
1013 path: Option<crate::analysis::table_sample::DrawPath>,
1014 ) -> Result<Self> {
1015 let mut view = source.sample_view(DataFrame::empty_with_schema(schema))?;
1016 let frame = scanned_frame(&view.original_lf)
1017 .ok_or_else(|| color_eyre::eyre::eyre!("a sample's frame has no rows to scan"))?;
1018 view.sampled = Some(Box::new(Sampled {
1019 source: Box::new(source),
1020 sample,
1021 rows,
1022 frame,
1023 through,
1024 drawn: None,
1025 path,
1026 }));
1027 Ok(view)
1028 }
1029
1030 pub fn sampled(&self) -> Option<&Sampled> {
1032 self.sampled.as_deref()
1033 }
1034
1035 pub fn unsampled(&self) -> &DataTableState {
1038 self.sampled
1039 .as_ref()
1040 .map_or(self, |sampled| sampled.source.as_ref())
1041 }
1042
1043 pub(crate) fn into_unsampled(mut self) -> DataTableState {
1046 match self.sampled.take() {
1047 Some(sampled) => *sampled.source,
1048 None => self,
1049 }
1050 }
1051
1052 pub(crate) fn sample_grew(&mut self) -> Option<bool> {
1056 let sampled = self.sampled.as_ref()?;
1057 let chunks = sampled.rows.take_new();
1058 if chunks.is_empty() {
1059 return None;
1060 }
1061 let mut frame = (*sampled.frame).clone();
1063 for chunk in &chunks {
1064 frame.vstack_mut(chunk).ok()?;
1065 }
1066 Some(self.rebind_sample(Arc::new(frame), false))
1067 }
1068
1069 pub(crate) fn sample_drawn(&mut self, drawn: crate::analysis::table_sample::Drawn) {
1072 let Some(sampled) = self.sampled.as_mut() else {
1073 return;
1074 };
1075 let ordered = sampled.rows.take_in_source_order().ok().flatten();
1078 sampled.path = drawn.path;
1080 sampled.drawn = Some(drawn);
1081 if let Some(frame) = ordered {
1082 self.rebind_sample(Arc::new(frame), true);
1083 }
1084 }
1085
1086 fn rebind_sample(&mut self, frame: Arc<DataFrame>, reordered: bool) -> bool {
1089 let Some(old) = self.sampled.as_ref().map(|sampled| sampled.frame.clone()) else {
1090 return false;
1091 };
1092 let rows_stand = !reordered
1093 && self.view.sort_columns.is_empty()
1094 && self.view.sort_ascending
1095 && self.scan_is_the_root();
1096 let rows = frame.height();
1097 self.each_frame(|lf| {
1098 crate::analysis::table_sample::rebind(&mut lf.logical_plan, &old, &frame)
1099 });
1100 if let Some(sampled) = self.sampled.as_mut() {
1101 sampled.frame = frame;
1102 }
1103 self.invalidate_num_rows();
1104 if self.is_pristine() {
1105 self.set_num_rows(rows);
1106 } else if self.scan_is_the_root() {
1107 self.pristine_rows = Some(rows);
1108 }
1109 if !rows_stand {
1110 self.drop_buffer();
1111 }
1112 self.needs_recollect = true;
1113 rows_stand
1114 }
1115
1116 pub(crate) fn sample_row_bytes(&self, from_source: bool) -> usize {
1119 let schema = if from_source {
1120 &self.original_schema
1121 } else {
1122 &self.view.schema
1123 };
1124 let columns: Vec<String> = schema
1125 .iter_names()
1126 .filter(|name| name.as_str() != crate::formats::schema_union::DRIFT_COLUMN)
1127 .map(|name| name.to_string())
1128 .collect();
1129 if !from_source && columns.len() == self.view.column_order.len() {
1131 return self.bytes_per_row();
1132 }
1133 estimate_bytes_per_row(schema, &columns, &self.column_bytes)
1134 }
1135
1136 pub fn files_a_page_reads(&self, start: usize, len: usize) -> Option<usize> {
1139 let offsets = self.files_window().and_then(|f| f.offsets.as_ref())?;
1140 let (first, last) = files_holding(offsets, start, len)?;
1141 Some(files_with_rows(offsets, first, last).len())
1142 }
1143
1144 pub(super) fn buffer_lf(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
1147 let mut all_columns = self.binary_stub_exprs();
1148 if self.carries_source_rows() {
1149 all_columns.push(col(crate::formats::schema_union::DRIFT_COLUMN));
1150 }
1151 self.window_lf(start, len, all_columns)
1152 }
1153
1154 pub(crate) fn carries_source_rows(&self) -> bool {
1157 self.view.drift_column_present
1158 || (self.scan_is_the_root() && (self.source_rows_at_open || self.view.view_numbered))
1159 }
1160
1161 pub fn row_numbers_from(&self, start: usize, rows: usize) -> Vec<usize> {
1164 let view = |i: usize| start + i + self.row_start_index;
1165 let places = self
1166 .view
1167 .buffered_df
1168 .as_ref()
1169 .filter(|_| self.carries_source_rows())
1170 .and_then(|df| df.column(crate::formats::schema_union::DRIFT_COLUMN).ok())
1171 .and_then(|column| {
1172 let offset = start.checked_sub(self.view.buffered_start_row)?;
1173 let len = rows.min(column.len().saturating_sub(offset));
1174 let slice = column.slice(offset as i64, len);
1175 let places = slice.u32().ok()?;
1176 let place = |p: usize| {
1178 self.numbering
1179 .as_ref()
1180 .and_then(|lines| lines.line_in_file(p))
1181 .unwrap_or(p)
1182 };
1183 Some(
1184 places
1185 .iter()
1186 .map(|p| p.map(|p| place(p as usize) + self.row_start_index))
1187 .collect::<Vec<_>>(),
1188 )
1189 });
1190 (0..rows)
1191 .map(|i| {
1192 places
1193 .as_ref()
1194 .and_then(|p| p.get(i).copied().flatten())
1195 .unwrap_or_else(|| view(i))
1196 })
1197 .collect()
1198 }
1199
1200 pub(super) fn window_lf(
1203 &self,
1204 start: usize,
1205 len: usize,
1206 all_columns: Vec<Expr>,
1207 ) -> PolarsResult<LazyFrame> {
1208 let records = self.window_now();
1209 window_of(
1210 &self.view.lf,
1211 Sources {
1212 files: self.files_window(),
1213 records: records.as_deref(),
1214 csv: self.csv_window(),
1215 },
1216 &self.read_as_text,
1217 start,
1218 len,
1219 all_columns,
1220 )
1221 }
1222
1223 pub(crate) fn view_rows(&self) -> ViewRows {
1226 ViewRows {
1227 lf: self.view.lf.clone(),
1228 files: self.files_window().cloned(),
1229 records: self.window_now().filter(|_| self.indexing().is_none()),
1232 csv: self.csv_window().cloned(),
1233 read_as_text: self.read_as_text.clone(),
1234 buffer: self
1235 .view
1236 .buffered_df
1237 .as_ref()
1238 .filter(|_| self.buffer_on_hand())
1239 .map(|df| (df.clone(), self.view.buffered_start_row)),
1240 num_rows: self.view.num_rows_valid.then_some(self.view.num_rows),
1241 streaming: self.polars_streaming,
1242 whole: sees_every_row_first(&self.view.lf),
1243 reads_up_to: reads_up_to_a_window(&self.view.lf),
1244 }
1245 }
1246
1247 pub(crate) fn go_to_found_row(&mut self, row: usize) -> bool {
1250 if !self.view.num_rows_valid && self.view.num_rows <= row {
1251 self.view.num_rows = row + 1;
1252 }
1253 self.scroll_to_row_centered(row)
1254 }
1255
1256 pub(crate) fn cursor_row(&self) -> usize {
1258 self.view.start_row + self.table_state.selected().unwrap_or(0)
1259 }
1260
1261 pub(super) fn bytes_per_row(&self) -> usize {
1264 self.view.observed_bytes_per_row.unwrap_or_else(|| {
1265 estimate_bytes_per_row(
1266 &self.view.schema,
1267 &self.view.column_order,
1268 &self.column_bytes,
1269 )
1270 })
1271 }
1272
1273 pub fn estimated_row_bytes(&self) -> usize {
1276 self.bytes_per_row()
1277 }
1278
1279 pub fn source_file_count(&self) -> Option<usize> {
1282 self.is_pristine().then(|| self.loaded_file_count())
1283 }
1284
1285 pub(crate) fn loaded_file_count(&self) -> usize {
1287 if !self.drift_files.is_empty() {
1288 self.drift_files.len()
1289 } else if let Some(remote) = &self.remote_files {
1290 remote.urls.len()
1291 } else {
1292 1
1293 }
1294 }
1295
1296 pub(super) fn byte_cap_rows(&self) -> usize {
1299 if self.max_buffered_mb == 0 {
1300 return 0;
1301 }
1302 let max_bytes = self.max_buffered_mb * 1024 * 1024;
1303 (max_bytes / self.bytes_per_row()).max(self.visible_rows.max(1))
1304 }
1305
1306 pub(super) fn remote_window(&self) -> bool {
1310 self.remote_source && self.is_pristine()
1311 }
1312
1313 fn csv_window(&self) -> Option<&Arc<super::csv_marks::CsvMarks>> {
1316 if !self.is_pristine() {
1317 return None;
1318 }
1319 self.csv_marks
1320 .get_or_init(|| super::csv_marks::CsvMarks::of(&self.original_lf))
1321 .as_ref()
1322 }
1323
1324 fn files_window(&self) -> Option<&RemoteFiles> {
1327 self.remote_files.as_ref().filter(|_| self.is_pristine())
1328 }
1329
1330 fn reach_rows(&self, pages: usize) -> usize {
1335 if !self.remote_window() || self.remote_files.is_some() {
1336 return pages * self.visible_rows.max(1);
1337 }
1338 let window = if self.max_buffered_rows > 0 {
1339 self.max_buffered_rows
1340 } else {
1341 DEFAULT_MAX_BUFFERED_ROWS
1342 };
1343 window / 2
1344 }
1345
1346 fn min_buffer_len(&self) -> usize {
1348 self.visible_rows.max(1) + self.reach_ahead() + self.reach_behind()
1349 }
1350
1351 pub fn at_end(&self) -> bool {
1353 self.view.start_row == self.view.num_rows.saturating_sub(self.visible_rows)
1354 }
1355
1356 fn fit_window(
1359 &self,
1360 view_start: usize,
1361 view_end: usize,
1362 buffer_start: &mut usize,
1363 buffer_end: &mut usize,
1364 ) {
1365 let byte_cap = self.byte_cap_rows();
1366 let cap = match (self.max_buffered_rows, byte_cap) {
1367 (0, cap) | (cap, 0) => cap,
1368 (rows, bytes) => rows.min(bytes),
1369 };
1370 if cap > 0 {
1371 shrink_around_view(
1372 view_start,
1373 view_end,
1374 cap,
1375 0,
1376 self.num_rows_bound(),
1377 buffer_start,
1378 buffer_end,
1379 );
1380 }
1381 let Some(offsets) = self
1382 .row_group_offsets
1383 .as_deref()
1384 .filter(|_| self.remote_window())
1385 else {
1386 if !self.remote_source {
1387 self.read_past_the_rows_on_hand(buffer_start, buffer_end);
1388 }
1389 return;
1390 };
1391 (*buffer_start, *buffer_end) = align_to_row_groups(
1392 offsets,
1393 view_start,
1394 view_end,
1395 *buffer_start,
1396 *buffer_end,
1397 cap,
1398 );
1399 if cap > 0 {
1402 let (floor, ceil) = (*buffer_start, *buffer_end);
1403 shrink_around_view(
1404 view_start,
1405 view_end,
1406 cap,
1407 floor,
1408 ceil,
1409 buffer_start,
1410 buffer_end,
1411 );
1412 }
1413 if let Some(file_offsets) = self.remote_files.as_ref().and_then(|f| f.offsets.as_ref()) {
1415 (*buffer_start, *buffer_end) = limit_files(
1416 file_offsets,
1417 view_start,
1418 view_end,
1419 *buffer_start,
1420 *buffer_end,
1421 MAX_FILES_PER_BUFFER,
1422 );
1423 }
1424 self.read_past_the_rows_on_hand(buffer_start, buffer_end);
1427 }
1428
1429 fn read_past_the_rows_on_hand(&self, buffer_start: &mut usize, buffer_end: &mut usize) {
1432 if !self.stitches_buffer() {
1433 return;
1434 }
1435 let (held_start, held_end) = (self.view.buffered_start_row, self.view.buffered_end_row);
1436 if held_start <= *buffer_start && *buffer_start < held_end && held_end < *buffer_end {
1437 *buffer_start = held_end;
1438 } else if *buffer_start < held_start && held_start < *buffer_end && *buffer_end <= held_end
1439 {
1440 *buffer_end = held_start;
1441 }
1442 }
1443
1444 fn release_display_buffer(&mut self) {
1447 self.widths.rows_arrived();
1448 self.view.buffered_df = None;
1449 self.view.locked_df = None;
1450 self.view.df = None;
1451 }
1452
1453 pub(super) fn slice_buffer_into_display(&mut self) {
1455 let full_df = match self.view.buffered_df.as_ref() {
1456 Some(df) => df,
1457 None => return,
1458 };
1459
1460 if self.view.locked_columns_count > 0 {
1461 let locked_names: Vec<&str> = self
1462 .view
1463 .column_order
1464 .iter()
1465 .take(self.view.locked_columns_count)
1466 .map(|s| s.as_str())
1467 .collect();
1468 if let Ok(locked_df) = full_df.select(locked_names) {
1469 self.view.locked_df = Some(locked_df);
1470 }
1471 } else {
1472 self.view.locked_df = None;
1473 }
1474
1475 let scroll_names: Vec<&str> = self
1476 .view
1477 .column_order
1478 .iter()
1479 .skip(self.frozen_shown() + self.termcol_index)
1480 .map(|s| s.as_str())
1481 .collect();
1482 if scroll_names.is_empty() {
1483 self.view.df = None;
1484 } else {
1485 if let Ok(scroll_df) = full_df.select(scroll_names) {
1486 self.view.df = Some(scroll_df);
1487 }
1488 }
1489 }
1490
1491 pub fn wants_to_load_ahead(&self) -> bool {
1494 if self.visible_rows == 0
1495 || self.view.buffered_df.is_none()
1496 || !self.page_on_hand(self.view.start_row)
1497 {
1498 return false;
1499 }
1500 let near = self.proximity();
1501 let view_end = self.view.start_row
1502 + self
1503 .visible_rows
1504 .min(self.num_rows_bound().saturating_sub(self.view.start_row));
1505 let behind = self.view.start_row - self.view.buffered_start_row <= self.proximity_behind()
1506 && self.view.buffered_start_row > 0;
1507 let ahead = self.view.buffered_end_row - view_end <= near
1508 && self.view.buffered_end_row < self.num_rows_bound();
1509 behind || ahead
1510 }
1511
1512 fn proximity(&self) -> usize {
1516 (self.reach_rows(self.pages_lookahead) / 2).max(self.visible_rows)
1517 }
1518
1519 fn proximity_behind(&self) -> usize {
1521 (self.reach_behind() / 2).max(self.visible_rows)
1522 }
1523
1524 fn reach_ahead(&self) -> usize {
1528 let reach = self.reach_rows(self.pages_lookahead);
1529 let measured = self
1530 .view
1531 .observed_bytes_per_row
1532 .filter(|_| self.decodes_pages && !self.remote_source && self.is_pristine());
1533 match measured {
1534 Some(bytes) => reach.max((BYTES_DECODED_AHEAD / bytes).min(PAGES_DECODED_AHEAD)),
1535 None => reach,
1536 }
1537 }
1538
1539 fn reach_behind(&self) -> usize {
1541 self.reach_rows(self.pages_lookback)
1542 }
1543
1544 pub fn buffer_position(&self) -> (u64, usize, usize, usize) {
1546 (
1547 self.len_generation(),
1548 self.view.start_row,
1549 self.view.buffered_start_row,
1550 self.view.buffered_end_row,
1551 )
1552 }
1553
1554 pub(crate) fn page_on_hand(&self, start: usize) -> bool {
1556 let bound = self.num_rows_bound();
1557 let end = start + self.visible_rows.min(bound.saturating_sub(start));
1558 self.view.buffered_df.is_some()
1559 && self.view.buffered_end_row > 0
1560 && start >= self.view.buffered_start_row
1561 && end <= self.view.buffered_end_row
1562 }
1563
1564 pub(crate) fn start_to_draw(&mut self) -> usize {
1567 if self.page_on_hand(self.view.start_row) {
1568 self.view.drawn_start = self.view.start_row;
1569 self.view.start_row
1570 } else if self.page_on_hand(self.view.drawn_start) {
1571 self.view.drawn_start
1572 } else {
1573 self.view.start_row
1574 }
1575 }
1576}