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 CollectResult {
68 df,
69 start,
70 returned,
71 bytes_per_row,
72 buffer_start: self.buffer_start,
73 buffer_end: self.buffer_end,
74 num_rows: self.num_rows,
75 count_known: self.count_known,
76 indexing: self.indexing,
77 }
78 }
79
80 fn cut_to_caps(&self, df: DataFrame, start: usize, seam: Option<usize>) -> (DataFrame, usize) {
85 let total = df.height();
86 if total == 0 {
87 return (df, start);
88 }
89 let mut max_rows = total;
91 if self.max_rows > 0 {
92 max_rows = max_rows.min(self.max_rows);
93 }
94 if self.max_mb > 0 {
95 let bytes_per_row = (df.estimated_size() / total).max(1);
96 max_rows = max_rows.min(self.max_mb * 1024 * 1024 / bytes_per_row);
97 }
98 let max_rows = max_rows.max(1);
99 if max_rows >= total {
100 return (df, start);
101 }
102 let view_off = self.view_start.saturating_sub(start).min(total);
103 let view_len = self.view_len.max(1).min(total);
104 let view_center = view_off + view_len / 2;
105 let mut keep_start = view_center.saturating_sub(max_rows / 2);
106 if keep_start + max_rows > total {
107 keep_start = total - max_rows;
108 }
109 let kept = max_rows.min(total - keep_start);
110 (trim_rows(df, keep_start, kept, seam), start + keep_start)
111 }
112}
113
114pub struct CollectResult {
117 pub(super) df: DataFrame,
119 pub(super) start: usize,
121 returned: usize,
123 bytes_per_row: Option<usize>,
125 buffer_start: usize,
127 buffer_end: usize,
128 num_rows: usize,
130 count_known: bool,
133 indexing: bool,
135}
136
137impl CollectResult {
138 pub(crate) fn rows(&self) -> &DataFrame {
140 &self.df
141 }
142}
143
144pub(super) const STRING_BYTES_GUESS: usize = 40;
146
147pub(super) fn estimate_bytes_per_row(
151 schema: &Schema,
152 columns: &[String],
153 column_bytes: &[(String, usize)],
154) -> usize {
155 let footer_width = |name: &String| {
156 column_bytes
157 .iter()
158 .find(|(n, _)| n == name)
159 .map(|(_, w)| *w)
160 };
161 columns
162 .iter()
163 .map(|name| match schema.get(name.as_str()) {
164 Some(DataType::String) => 16 + footer_width(name).unwrap_or(STRING_BYTES_GUESS - 16),
165 Some(DataType::Binary) => 16 + binary_stub().len(),
166 Some(DataType::Boolean) => 1,
167 Some(DataType::Null) => 0,
168 Some(dtype) if dtype.is_primitive_numeric() || dtype.is_temporal() => {
169 match dtype.to_physical() {
170 DataType::Int8 | DataType::UInt8 => 1,
171 DataType::Int16 | DataType::UInt16 => 2,
172 DataType::Int32 | DataType::UInt32 | DataType::Float32 => 4,
173 DataType::Int128 => 16,
174 _ => 8,
175 }
176 }
177 Some(DataType::Decimal(..)) => 16,
178 _ => footer_width(name).unwrap_or(64),
179 })
180 .sum::<usize>()
181 .max(1)
182}
183
184pub(super) fn trim_rows(
189 df: DataFrame,
190 offset: usize,
191 len: usize,
192 seam: Option<usize>,
193) -> DataFrame {
194 if backing_rows(&df, offset, len) > len + len / 4 {
195 compact_rows(df, offset, len, seam)
196 } else {
197 df.slice(offset as i64, len)
198 }
199}
200
201pub(super) fn backing_rows(df: &DataFrame, offset: usize, len: usize) -> usize {
204 let end = offset + len;
205 df.columns()
206 .iter()
207 .filter_map(Column::as_series)
208 .map(|s| {
209 let mut start = 0;
210 let mut touched = 0;
211 for chunk in s.chunks() {
212 let chunk_end = start + chunk.len();
213 if start < end && offset < chunk_end {
214 touched += chunk.len();
215 }
216 start = chunk_end;
217 }
218 touched
219 })
220 .max()
221 .unwrap_or(len)
222}
223
224pub(super) fn compact_rows(
230 df: DataFrame,
231 offset: usize,
232 len: usize,
233 seam: Option<usize>,
234) -> DataFrame {
235 use polars::series::builder::SeriesBuilder;
236 use polars_arrow::array::builder::ShareStrategy;
237 #[cfg(test)]
238 tests::COMPACTIONS.with(|count| count.set(count.get() + 1));
239 let len = len.min(df.height().saturating_sub(offset));
240 let pieces = match seam.filter(|&seam| offset < seam && seam < offset + len) {
241 Some(seam) => vec![(offset, seam - offset), (seam, offset + len - seam)],
242 None => vec![(offset, len)],
243 };
244 let copy = |series: &Series, (offset, len): (usize, usize)| {
245 let mut builder = SeriesBuilder::new(series.dtype().clone());
246 builder.reserve(len);
247 builder.subslice_extend(series, offset, len, ShareStrategy::Never);
248 builder.freeze(series.name().clone())
249 };
250 let columns = df
251 .into_columns()
252 .into_iter()
253 .map(|column| match column {
254 Column::Scalar(constant) => {
255 Column::new_scalar(constant.name().clone(), constant.scalar().clone(), len)
256 }
257 Column::Series(series) => {
258 let mut kept = copy(&series, pieces[0]);
259 for &piece in &pieces[1..] {
260 if kept.append_owned(copy(&series, piece)).is_err() {
261 kept = copy(&series, (offset, len));
262 break;
263 }
264 }
265 kept.into_column()
266 }
267 })
268 .collect();
269 DataFrame::new(len, columns).unwrap_or_else(|_| DataFrame::empty_with_height(len))
271}
272
273pub(super) fn shrink_around_view(
276 view_start: usize,
277 view_end: usize,
278 max_len: usize,
279 floor: usize,
280 ceil: usize,
281 buffer_start: &mut usize,
282 buffer_end: &mut usize,
283) {
284 if buffer_end.saturating_sub(*buffer_start) <= max_len {
285 return;
286 }
287 let view_len = view_end.saturating_sub(view_start);
288 if view_len >= max_len {
289 *buffer_start = view_start;
290 *buffer_end = (view_start + max_len).min(ceil);
291 return;
292 }
293 let half = (max_len - view_len) / 2;
294 *buffer_end = (view_end + half).min(ceil);
295 *buffer_start = buffer_end.saturating_sub(max_len).max(floor);
296 if *buffer_start > view_start {
297 *buffer_start = view_start;
298 }
299 *buffer_end = (*buffer_start + max_len).min(ceil);
300}
301
302pub(super) const MAX_FILES_PER_BUFFER: usize = 16;
304
305pub(super) fn limit_files(
308 offsets: &[usize],
309 view_start: usize,
310 view_end: usize,
311 start: usize,
312 end: usize,
313 max_files: usize,
314) -> (usize, usize) {
315 let (Some((first, last)), Some((view_first, view_last))) = (
316 files_holding(offsets, start, end.saturating_sub(start)),
317 files_holding(
318 offsets,
319 view_start,
320 view_end.saturating_sub(view_start).max(1),
321 ),
322 ) else {
323 return (start, end);
324 };
325 let opened = |from: usize, to: usize| (from..=to).filter(|&i| holds_rows(offsets, i)).count();
327 if opened(first, last) <= max_files {
328 return (start, end);
329 }
330 let (mut lo, mut hi) = (view_first.max(first), view_last.min(last));
331 let mut files = opened(lo, hi);
332 while files < max_files && (hi < last || lo > first) {
333 if hi < last {
334 hi += 1;
335 files += usize::from(holds_rows(offsets, hi));
336 }
337 if files < max_files && lo > first {
338 lo -= 1;
339 files += usize::from(holds_rows(offsets, lo));
340 }
341 }
342 (start.max(offsets[lo]), end.min(offsets[hi + 1]))
343}
344
345pub(super) fn holds_rows(offsets: &[usize], i: usize) -> bool {
347 offsets[i + 1] > offsets[i]
348}
349
350pub(super) fn files_with_rows(offsets: &[usize], first: usize, last: usize) -> Vec<usize> {
353 (first..=last).filter(|&i| holds_rows(offsets, i)).collect()
354}
355
356pub(super) fn files_holding(offsets: &[usize], start: usize, len: usize) -> Option<(usize, usize)> {
359 let files = offsets.len().checked_sub(1)?;
360 let total = *offsets.last()?;
361 if files == 0 || len == 0 || start >= total {
362 return None;
363 }
364 let end = (start + len).min(total);
365 let file_of = |row: usize| offsets.partition_point(|&o| o <= row).saturating_sub(1);
368 Some((file_of(start), file_of(end - 1).min(files - 1)))
369}
370
371pub(super) fn window_of(
374 lf: &LazyFrame,
375 files: Option<&RemoteFiles>,
376 records: Option<&dyn crate::formats::pushdown::Windowed>,
377 read_as_text: &[PlSmallStr],
378 start: usize,
379 len: usize,
380 all_columns: Vec<Expr>,
381) -> PolarsResult<LazyFrame> {
382 if let Some(records) = records {
385 return Ok(records.window(start, len)?.select(all_columns));
386 }
387 if let Some((files, offsets)) = files.and_then(|f| f.offsets.as_ref().map(|o| (f, o)))
388 && let Some((first, last)) = files_holding(offsets, start, len)
389 {
390 let urls: Vec<String> = files_with_rows(offsets, first, last)
393 .into_iter()
394 .map(|i| files.urls[i].clone())
395 .collect();
396 let lf = (files.scan)(&urls, read_as_text)?;
397 return Ok(lf
398 .select(all_columns)
399 .slice((start - offsets[first]) as i64, len as u32));
400 }
401 Ok(lf
402 .clone()
403 .select(all_columns)
404 .slice(start as i64, len as u32))
405}
406
407#[derive(Clone)]
410pub(crate) struct ViewRows {
411 lf: LazyFrame,
412 files: Option<RemoteFiles>,
413 records: Option<Arc<dyn crate::formats::pushdown::Windowed>>,
415 read_as_text: Vec<PlSmallStr>,
416 pub(crate) buffer: Option<(DataFrame, usize)>,
418 pub(crate) num_rows: Option<usize>,
420 pub(crate) streaming: bool,
421 pub(crate) whole: bool,
423 pub(crate) reads_up_to: bool,
425}
426
427pub(crate) fn sees_every_row_first(lf: &LazyFrame) -> bool {
430 use polars::lazy::dsl::DslPlan;
431 lf.logical_plan.into_iter().any(|node| {
432 matches!(
433 node,
434 DslPlan::Sort { .. } | DslPlan::GroupBy { .. } | DslPlan::Pivot { .. }
435 )
436 })
437}
438
439pub(crate) fn reads_up_to_a_window(lf: &LazyFrame) -> bool {
442 use polars::lazy::dsl::{DslPlan, FileScanDsl};
443 lf.logical_plan.into_iter().any(|node| match node {
444 DslPlan::Filter { .. } => true,
445 DslPlan::Scan { scan_type, .. } => !matches!(
446 **scan_type,
447 FileScanDsl::Parquet { .. } | FileScanDsl::Ipc { .. }
448 ),
449 _ => false,
450 })
451}
452
453impl ViewRows {
454 pub(crate) fn window(
456 &self,
457 start: usize,
458 len: usize,
459 exprs: Vec<Expr>,
460 ) -> PolarsResult<LazyFrame> {
461 window_of(
462 &self.lf,
463 self.files.as_ref(),
464 self.records.as_deref(),
465 &self.read_as_text,
466 start,
467 len,
468 exprs,
469 )
470 }
471
472 #[cfg(test)]
474 pub(crate) fn of(lf: LazyFrame, buffer: Option<(DataFrame, usize)>) -> Self {
475 Self {
476 whole: sees_every_row_first(&lf),
477 reads_up_to: reads_up_to_a_window(&lf),
478 lf,
479 files: None,
480 records: None,
481 read_as_text: Vec::new(),
482 buffer,
483 num_rows: None,
484 streaming: false,
485 }
486 }
487}
488
489pub(super) fn align_to_row_groups(
494 offsets: &[usize],
495 view_start: usize,
496 view_end: usize,
497 start: usize,
498 end: usize,
499 cap: usize,
500) -> (usize, usize) {
501 let Some(groups) = offsets.len().checked_sub(1).filter(|n| *n > 0) else {
502 return (start, end);
503 };
504 let group_of = |row: usize| {
505 offsets
506 .partition_point(|&o| o <= row)
507 .saturating_sub(1)
508 .min(groups - 1)
509 };
510 let last_row = |s: usize, e: usize| e.saturating_sub(1).max(s);
511 let (mut lo, mut hi) = (
512 group_of(view_start),
513 group_of(last_row(view_start, view_end)),
514 );
515 let (want_lo, want_hi) = (group_of(start), group_of(last_row(start, end)));
516 let fits = |lo: usize, hi: usize| cap == 0 || offsets[hi + 1] - offsets[lo] <= cap;
517 loop {
518 if hi < want_hi && fits(lo, hi + 1) {
519 hi += 1;
520 } else if lo > want_lo && fits(lo - 1, hi) {
521 lo -= 1;
522 } else {
523 break;
524 }
525 }
526 (offsets[lo], offsets[hi + 1])
527}
528
529impl DataTableState {
530 pub fn scroll_would_trigger_collect(&self, rows: i64) -> bool {
533 if rows < 0 && self.view.start_row == 0 {
534 return false;
535 }
536 let new_start_row = if self.view.start_row as i64 + rows <= 0 {
537 0
538 } else {
539 if let Some(df) = self.view.df.as_ref()
540 && rows > 0
541 && df.shape().0 <= self.visible_rows
542 {
543 return false;
544 }
545 let unclamped = (self.view.start_row as i64 + rows) as usize;
546 if rows > 0 {
547 unclamped.min(self.view.num_rows.saturating_sub(self.visible_rows))
548 } else {
549 unclamped
550 }
551 };
552 if new_start_row == self.view.start_row {
553 return false;
554 }
555 let view_end = new_start_row
556 + self
557 .visible_rows
558 .min(self.view.num_rows.saturating_sub(new_start_row));
559 let within_buffer = new_start_row >= self.view.buffered_start_row
560 && view_end <= self.view.buffered_end_row
561 && self.view.buffered_end_row > 0;
562 !within_buffer
563 }
564
565 pub fn slide_table(&mut self, rows: i64) -> bool {
568 if rows < 0 && self.view.start_row == 0 {
569 return false;
570 }
571
572 let new_start_row = if self.view.start_row as i64 + rows <= 0 {
573 0
574 } else {
575 if let Some(df) = self.view.df.as_ref()
576 && rows > 0
577 && df.shape().0 <= self.visible_rows
578 {
579 return false;
580 }
581 let unclamped = (self.view.start_row as i64 + rows) as usize;
582 if rows > 0 {
583 unclamped.min(self.view.num_rows.saturating_sub(self.visible_rows))
586 } else {
587 unclamped
588 }
589 };
590
591 if new_start_row == self.view.start_row {
592 return false;
593 }
594
595 let view_end = new_start_row
596 + self
597 .visible_rows
598 .min(self.view.num_rows.saturating_sub(new_start_row));
599 let within_buffer = new_start_row >= self.view.buffered_start_row
600 && view_end <= self.view.buffered_end_row
601 && self.view.buffered_end_row > 0;
602
603 self.view.start_row = new_start_row;
604
605 if within_buffer {
606 if self.table_state.selected().is_none() {
607 self.table_state.select(Some(0));
608 }
609 false
610 } else {
611 true }
613 }
614
615 #[cfg(test)]
617 pub fn collect(&mut self) {
618 if self.defer_collect {
619 return;
620 }
621 if !self.view.num_rows_valid {
622 match collect_lazy(row_count_lf(&self.view.lf), self.polars_streaming) {
625 Ok(df) => {
626 self.error = None;
627 let n = match df.get(0).as_deref().and_then(|row| row.first()) {
628 Some(AnyValue::UInt64(len)) => *len as usize,
629 _ => 0,
630 };
631 self.set_num_rows(n);
632 }
633 Err(e) => {
634 self.error = Some(e);
635 self.set_num_rows(0);
636 }
637 }
638 }
639 let Some(request) = self.prepare_async_collect(None) else {
640 return;
641 };
642 match collect_lazy(request.lf, request.polars_streaming) {
643 Ok(df) => self.apply_async_collect(request.plan.fit(df)),
644 Err(e) => self.error = Some(e),
645 }
646 }
647
648 #[cfg(not(test))]
651 pub(super) fn collect(&mut self) {
652 if !self.defer_collect {
653 self.needs_recollect = true;
654 }
655 }
656
657 pub(crate) fn binary_stub_exprs(&self) -> Vec<Expr> {
661 self.view
662 .column_order
663 .iter()
664 .map(|name| {
665 if matches!(self.view.schema.get(name.as_str()), Some(DataType::Binary)) {
666 lit(binary_stub()).alias(name.as_str())
667 } else {
668 col(name.as_str())
669 }
670 })
671 .collect()
672 }
673
674 pub fn prepare_async_collect(
678 &mut self,
679 num_rows_override: Option<usize>,
680 ) -> Option<CollectRequest> {
681 if self.visible_rows > 0 {
682 self.proximity_threshold = self.proximity();
683 }
684
685 if let Some(n) = num_rows_override {
686 self.view.num_rows = n;
687 self.view.num_rows_valid = true;
688 }
689
690 let count_known = self.view.num_rows_valid;
693 let bound = self.num_rows_bound();
694
695 if count_known {
696 if self.view.num_rows > 0 {
697 let max_start = self.view.num_rows.saturating_sub(1);
698 if self.view.start_row > max_start {
699 self.view.start_row = max_start;
700 }
701 } else {
702 self.view.start_row = 0;
704 self.drop_buffer();
705 self.view.df = None;
706 self.view.locked_df = None;
707 return None;
708 }
709 }
710
711 if self.view.column_order.is_empty() {
714 self.drop_buffer();
715 self.view.df = None;
716 self.view.locked_df = None;
717 return None;
718 }
719
720 let view_start = self.view.start_row;
721 let view_end = self.view.start_row + self.visible_rows.min(bound - self.view.start_row);
722 let within_buffer = view_start >= self.view.buffered_start_row
723 && view_end <= self.view.buffered_end_row
724 && self.view.buffered_end_row > 0;
725
726 let (new_buffer_start, new_buffer_end) = if within_buffer {
728 let dist_to_start = view_start.saturating_sub(self.view.buffered_start_row);
729 let dist_to_end = self.view.buffered_end_row.saturating_sub(view_end);
730 let needs_expansion_back =
731 dist_to_start <= self.proximity_threshold && self.view.buffered_start_row > 0;
732 let needs_expansion_forward =
733 dist_to_end <= self.proximity_threshold && self.view.buffered_end_row < bound;
734
735 if !needs_expansion_back && !needs_expansion_forward {
736 (self.view.buffered_start_row, self.view.buffered_end_row)
738 } else {
739 let mut s = if needs_expansion_back {
740 view_start.saturating_sub(self.reach_rows(self.pages_lookback))
741 } else {
742 self.view.buffered_start_row
743 };
744 let mut e = if needs_expansion_forward {
745 (view_end + self.reach_rows(self.pages_lookahead)).min(bound)
746 } else {
747 self.view.buffered_end_row
748 };
749 self.fit_window(view_start, view_end, &mut s, &mut e);
750 (s, e)
751 }
752 } else {
753 let had_buffer = self.view.buffered_end_row > 0;
754 let scrolled_past_end = had_buffer && view_start >= self.view.buffered_end_row;
755 let scrolled_past_start = had_buffer && view_end <= self.view.buffered_start_row;
756 let extend_forward_ok = scrolled_past_end
757 && (view_start - self.view.buffered_end_row)
758 <= self.reach_rows(self.pages_lookahead);
759 let extend_backward_ok = scrolled_past_start
760 && (self.view.buffered_start_row - view_end)
761 <= self.reach_rows(self.pages_lookback);
762
763 let mut s;
764 let mut e;
765 if extend_forward_ok {
766 s = self.view.buffered_start_row;
767 e = (view_end + self.reach_rows(self.pages_lookahead)).min(bound);
768 } else if extend_backward_ok {
769 s = view_start.saturating_sub(self.reach_rows(self.pages_lookback));
770 e = self.view.buffered_end_row;
771 } else {
772 s = view_start.saturating_sub(self.reach_rows(self.pages_lookback));
773 e = (view_end + self.reach_rows(self.pages_lookahead)).min(bound);
774 let min_initial_len = self.min_buffer_len();
775 let current_len = e.saturating_sub(s);
776 if current_len < min_initial_len {
777 let need = min_initial_len.saturating_sub(current_len);
778 let can_extend_end = bound.saturating_sub(e);
779 let can_extend_start = s;
780 if can_extend_end >= need {
781 e = (e + need).min(bound);
782 } else if can_extend_start >= need {
783 s = s.saturating_sub(need);
784 } else {
785 e = (e + can_extend_end).min(bound);
786 s = s.saturating_sub(need.saturating_sub(can_extend_end));
787 }
788 }
789 }
790 self.fit_window(view_start, view_end, &mut s, &mut e);
791 (s, e)
792 };
793
794 let buffer_size = new_buffer_end.saturating_sub(new_buffer_start);
795 if buffer_size == 0 {
796 return None;
797 }
798 if self.holds_buffer(new_buffer_start, new_buffer_end) {
801 self.slice_buffer_into_display();
802 if self.table_state.selected().is_none() {
803 self.table_state.select(Some(0));
804 }
805 return None;
806 }
807
808 let lf = match self.buffer_lf(new_buffer_start, buffer_size) {
809 Ok(lf) => lf,
810 Err(e) => {
811 self.error = Some(e);
812 return None;
813 }
814 };
815
816 let num_rows = if count_known {
819 self.view.num_rows
820 } else {
821 new_buffer_end
822 };
823 Some(CollectRequest {
824 lf,
825 polars_streaming: self.polars_streaming,
826 buffer_start: new_buffer_start,
827 buffer_end: new_buffer_end,
828 plan: self.fill_plan(new_buffer_start, new_buffer_end, num_rows, count_known),
829 })
830 }
831
832 pub(super) fn fill_plan(
835 &self,
836 buffer_start: usize,
837 buffer_end: usize,
838 num_rows: usize,
839 count_known: bool,
840 ) -> FillPlan {
841 let held = self
842 .abuts_buffer(buffer_start, buffer_end.saturating_sub(buffer_start))
843 .then(|| self.view.buffered_df.clone())
844 .flatten()
845 .map(|df| (df, self.view.buffered_start_row));
846 FillPlan {
847 buffer_start,
848 buffer_end,
849 num_rows,
850 count_known,
851 indexing: self.indexing().is_some(),
852 held,
853 view_start: self.view.start_row,
854 view_len: self.visible_rows,
855 max_rows: self.max_buffered_rows,
856 max_mb: self.max_buffered_mb,
857 }
858 }
859
860 pub fn apply_async_collect(&mut self, result: CollectResult) {
863 let CollectResult {
864 df,
865 start,
866 returned: returned_rows,
867 bytes_per_row,
868 buffer_start,
869 buffer_end,
870 num_rows,
871 count_known,
872 indexing,
873 } = result;
874 let requested_rows = buffer_end.saturating_sub(buffer_start);
875
876 if count_known {
877 self.view.num_rows = num_rows;
878 self.view.num_rows_valid = true;
879 } else if returned_rows < requested_rows
880 && (buffer_start == 0 || returned_rows > 0)
881 && !indexing
883 && self.indexing().is_none()
884 {
885 self.view.num_rows = buffer_start + returned_rows;
888 self.view.num_rows_valid = true;
889 } else if !self.view.num_rows_valid {
890 self.view.num_rows = self.view.num_rows.max(buffer_end);
893 }
894 self.error = None;
897 self.remember_pristine_count();
898
899 if bytes_per_row.is_some() {
900 self.view.observed_bytes_per_row = bytes_per_row;
901 }
902 let end = start + df.height();
906 let view_end = self.view.start_row + self.visible_rows.max(1);
907 let reaches_end = end >= buffer_start + returned_rows;
908 let shows_view = start <= self.view.start_row
909 && (self.view.start_row < end || (returned_rows < requested_rows && reaches_end));
910 if !shows_view {
911 self.needs_recollect = true;
912 return;
913 }
914 self.release_display_buffer();
915 self.view.buffered_start_row = start;
916 self.view.buffered_end_row = end;
917 self.view.buffered_df = Some(df);
918 self.slice_buffer_into_display();
920 if self.table_state.selected().is_none() {
921 self.table_state.select(Some(0));
922 }
923 if view_end > end && end < self.view.num_rows {
924 self.needs_recollect = true;
925 }
926 }
927
928 fn abuts_buffer(&self, start: usize, rows: usize) -> bool {
931 self.stitches_buffer()
932 && (start == self.view.buffered_end_row || start + rows == self.view.buffered_start_row)
933 }
934
935 pub(crate) fn sampled_from(
939 source: DataTableState,
940 sample: crate::analysis::sampling::Sample,
941 schema: &Schema,
942 rows: Arc<crate::analysis::table_sample::SampleRows>,
943 through: bool,
944 path: Option<crate::analysis::table_sample::DrawPath>,
945 ) -> Result<Self> {
946 let mut view = source.sample_view(DataFrame::empty_with_schema(schema))?;
947 let frame = scanned_frame(&view.original_lf)
948 .ok_or_else(|| color_eyre::eyre::eyre!("a sample's frame has no rows to scan"))?;
949 view.sampled = Some(Box::new(Sampled {
950 source: Box::new(source),
951 sample,
952 rows,
953 frame,
954 through,
955 drawn: None,
956 path,
957 }));
958 Ok(view)
959 }
960
961 pub fn sampled(&self) -> Option<&Sampled> {
963 self.sampled.as_deref()
964 }
965
966 pub fn unsampled(&self) -> &DataTableState {
969 self.sampled
970 .as_ref()
971 .map_or(self, |sampled| sampled.source.as_ref())
972 }
973
974 pub(crate) fn into_unsampled(mut self) -> DataTableState {
977 match self.sampled.take() {
978 Some(sampled) => *sampled.source,
979 None => self,
980 }
981 }
982
983 pub(crate) fn sample_grew(&mut self) -> Option<bool> {
987 let sampled = self.sampled.as_ref()?;
988 let chunks = sampled.rows.take_new();
989 if chunks.is_empty() {
990 return None;
991 }
992 let mut frame = (*sampled.frame).clone();
994 for chunk in &chunks {
995 frame.vstack_mut(chunk).ok()?;
996 }
997 Some(self.rebind_sample(Arc::new(frame), false))
998 }
999
1000 pub(crate) fn sample_drawn(&mut self, drawn: crate::analysis::table_sample::Drawn) {
1003 let Some(sampled) = self.sampled.as_mut() else {
1004 return;
1005 };
1006 let ordered = sampled.rows.take_in_source_order().ok().flatten();
1009 sampled.path = drawn.path;
1011 sampled.drawn = Some(drawn);
1012 if let Some(frame) = ordered {
1013 self.rebind_sample(Arc::new(frame), true);
1014 }
1015 }
1016
1017 fn rebind_sample(&mut self, frame: Arc<DataFrame>, reordered: bool) -> bool {
1020 let Some(old) = self.sampled.as_ref().map(|sampled| sampled.frame.clone()) else {
1021 return false;
1022 };
1023 let rows_stand = !reordered
1024 && self.view.sort_columns.is_empty()
1025 && self.view.sort_ascending
1026 && self.scan_is_the_root();
1027 let rows = frame.height();
1028 self.each_frame(|lf| {
1029 crate::analysis::table_sample::rebind(&mut lf.logical_plan, &old, &frame)
1030 });
1031 if let Some(sampled) = self.sampled.as_mut() {
1032 sampled.frame = frame;
1033 }
1034 self.invalidate_num_rows();
1035 if self.is_pristine() {
1036 self.set_num_rows(rows);
1037 } else if self.scan_is_the_root() {
1038 self.pristine_rows = Some(rows);
1039 }
1040 if !rows_stand {
1041 self.drop_buffer();
1042 }
1043 self.needs_recollect = true;
1044 rows_stand
1045 }
1046
1047 pub(crate) fn sample_row_bytes(&self, from_source: bool) -> usize {
1050 let schema = if from_source {
1051 &self.original_schema
1052 } else {
1053 &self.view.schema
1054 };
1055 let columns: Vec<String> = schema
1056 .iter_names()
1057 .filter(|name| name.as_str() != crate::formats::schema_union::DRIFT_COLUMN)
1058 .map(|name| name.to_string())
1059 .collect();
1060 if !from_source && columns.len() == self.view.column_order.len() {
1062 return self.bytes_per_row();
1063 }
1064 estimate_bytes_per_row(schema, &columns, &self.column_bytes)
1065 }
1066
1067 pub fn files_a_page_reads(&self, start: usize, len: usize) -> Option<usize> {
1070 let offsets = self.files_window().and_then(|f| f.offsets.as_ref())?;
1071 let (first, last) = files_holding(offsets, start, len)?;
1072 Some(files_with_rows(offsets, first, last).len())
1073 }
1074
1075 pub(super) fn buffer_lf(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
1078 let mut all_columns = self.binary_stub_exprs();
1079 if self.carries_source_rows() {
1080 all_columns.push(col(crate::formats::schema_union::DRIFT_COLUMN));
1081 }
1082 self.window_lf(start, len, all_columns)
1083 }
1084
1085 pub(crate) fn carries_source_rows(&self) -> bool {
1088 self.view.drift_column_present
1089 || (self.scan_is_the_root() && (self.source_rows_at_open || self.view.view_numbered))
1090 }
1091
1092 pub fn row_numbers_from(&self, start: usize, rows: usize) -> Vec<usize> {
1095 let view = |i: usize| start + i + self.row_start_index;
1096 let places = self
1097 .view
1098 .buffered_df
1099 .as_ref()
1100 .filter(|_| self.carries_source_rows())
1101 .and_then(|df| df.column(crate::formats::schema_union::DRIFT_COLUMN).ok())
1102 .and_then(|column| {
1103 let offset = start.checked_sub(self.view.buffered_start_row)?;
1104 let len = rows.min(column.len().saturating_sub(offset));
1105 let slice = column.slice(offset as i64, len);
1106 let places = slice.u32().ok()?;
1107 let place = |p: usize| {
1109 self.numbering
1110 .as_ref()
1111 .and_then(|lines| lines.line_in_file(p))
1112 .unwrap_or(p)
1113 };
1114 Some(
1115 places
1116 .iter()
1117 .map(|p| p.map(|p| place(p as usize) + self.row_start_index))
1118 .collect::<Vec<_>>(),
1119 )
1120 });
1121 (0..rows)
1122 .map(|i| {
1123 places
1124 .as_ref()
1125 .and_then(|p| p.get(i).copied().flatten())
1126 .unwrap_or_else(|| view(i))
1127 })
1128 .collect()
1129 }
1130
1131 pub(super) fn window_lf(
1134 &self,
1135 start: usize,
1136 len: usize,
1137 all_columns: Vec<Expr>,
1138 ) -> PolarsResult<LazyFrame> {
1139 window_of(
1140 &self.view.lf,
1141 self.files_window(),
1142 self.window_now().as_deref(),
1143 &self.read_as_text,
1144 start,
1145 len,
1146 all_columns,
1147 )
1148 }
1149
1150 pub(crate) fn view_rows(&self) -> ViewRows {
1153 ViewRows {
1154 lf: self.view.lf.clone(),
1155 files: self.files_window().cloned(),
1156 records: self.window_now().filter(|_| self.indexing().is_none()),
1159 read_as_text: self.read_as_text.clone(),
1160 buffer: self
1161 .view
1162 .buffered_df
1163 .as_ref()
1164 .filter(|_| self.buffer_on_hand())
1165 .map(|df| (df.clone(), self.view.buffered_start_row)),
1166 num_rows: self.view.num_rows_valid.then_some(self.view.num_rows),
1167 streaming: self.polars_streaming,
1168 whole: sees_every_row_first(&self.view.lf),
1169 reads_up_to: reads_up_to_a_window(&self.view.lf),
1170 }
1171 }
1172
1173 pub(crate) fn go_to_found_row(&mut self, row: usize) -> bool {
1176 if !self.view.num_rows_valid && self.view.num_rows <= row {
1177 self.view.num_rows = row + 1;
1178 }
1179 self.scroll_to_row_centered(row)
1180 }
1181
1182 pub(crate) fn cursor_row(&self) -> usize {
1184 self.view.start_row + self.table_state.selected().unwrap_or(0)
1185 }
1186
1187 pub(super) fn bytes_per_row(&self) -> usize {
1190 self.view.observed_bytes_per_row.unwrap_or_else(|| {
1191 estimate_bytes_per_row(
1192 &self.view.schema,
1193 &self.view.column_order,
1194 &self.column_bytes,
1195 )
1196 })
1197 }
1198
1199 pub fn estimated_row_bytes(&self) -> usize {
1202 self.bytes_per_row()
1203 }
1204
1205 pub fn source_file_count(&self) -> Option<usize> {
1208 self.is_pristine().then(|| self.loaded_file_count())
1209 }
1210
1211 pub(crate) fn loaded_file_count(&self) -> usize {
1213 if !self.drift_files.is_empty() {
1214 self.drift_files.len()
1215 } else if let Some(remote) = &self.remote_files {
1216 remote.urls.len()
1217 } else {
1218 1
1219 }
1220 }
1221
1222 pub(super) fn byte_cap_rows(&self) -> usize {
1225 if self.max_buffered_mb == 0 {
1226 return 0;
1227 }
1228 let max_bytes = self.max_buffered_mb * 1024 * 1024;
1229 (max_bytes / self.bytes_per_row()).max(self.visible_rows.max(1))
1230 }
1231
1232 pub(super) fn remote_window(&self) -> bool {
1236 self.remote_source && self.is_pristine()
1237 }
1238
1239 fn files_window(&self) -> Option<&RemoteFiles> {
1242 self.remote_files.as_ref().filter(|_| self.is_pristine())
1243 }
1244
1245 fn reach_rows(&self, pages: usize) -> usize {
1250 if !self.remote_window() || self.remote_files.is_some() {
1251 return pages * self.visible_rows.max(1);
1252 }
1253 let window = if self.max_buffered_rows > 0 {
1254 self.max_buffered_rows
1255 } else {
1256 DEFAULT_MAX_BUFFERED_ROWS
1257 };
1258 window / 2
1259 }
1260
1261 fn min_buffer_len(&self) -> usize {
1263 self.visible_rows.max(1)
1264 + self.reach_rows(self.pages_lookahead)
1265 + self.reach_rows(self.pages_lookback)
1266 }
1267
1268 pub fn at_end(&self) -> bool {
1270 self.view.start_row == self.view.num_rows.saturating_sub(self.visible_rows)
1271 }
1272
1273 fn fit_window(
1276 &self,
1277 view_start: usize,
1278 view_end: usize,
1279 buffer_start: &mut usize,
1280 buffer_end: &mut usize,
1281 ) {
1282 let byte_cap = self.byte_cap_rows();
1283 let cap = match (self.max_buffered_rows, byte_cap) {
1284 (0, cap) | (cap, 0) => cap,
1285 (rows, bytes) => rows.min(bytes),
1286 };
1287 if cap > 0 {
1288 shrink_around_view(
1289 view_start,
1290 view_end,
1291 cap,
1292 0,
1293 self.num_rows_bound(),
1294 buffer_start,
1295 buffer_end,
1296 );
1297 }
1298 let Some(offsets) = self
1299 .row_group_offsets
1300 .as_deref()
1301 .filter(|_| self.remote_window())
1302 else {
1303 return;
1304 };
1305 (*buffer_start, *buffer_end) = align_to_row_groups(
1306 offsets,
1307 view_start,
1308 view_end,
1309 *buffer_start,
1310 *buffer_end,
1311 cap,
1312 );
1313 if cap > 0 {
1316 let (floor, ceil) = (*buffer_start, *buffer_end);
1317 shrink_around_view(
1318 view_start,
1319 view_end,
1320 cap,
1321 floor,
1322 ceil,
1323 buffer_start,
1324 buffer_end,
1325 );
1326 }
1327 if let Some(file_offsets) = self.remote_files.as_ref().and_then(|f| f.offsets.as_ref()) {
1329 (*buffer_start, *buffer_end) = limit_files(
1330 file_offsets,
1331 view_start,
1332 view_end,
1333 *buffer_start,
1334 *buffer_end,
1335 MAX_FILES_PER_BUFFER,
1336 );
1337 }
1338 if self.buffer_on_hand() {
1341 let (held_start, held_end) = (self.view.buffered_start_row, self.view.buffered_end_row);
1342 if held_start <= *buffer_start && *buffer_start < held_end && held_end < *buffer_end {
1343 *buffer_start = held_end;
1344 } else if *buffer_start < held_start
1345 && held_start < *buffer_end
1346 && *buffer_end <= held_end
1347 {
1348 *buffer_end = held_start;
1349 }
1350 }
1351 }
1352
1353 fn release_display_buffer(&mut self) {
1356 self.widths.rows_arrived();
1357 self.view.buffered_df = None;
1358 self.view.locked_df = None;
1359 self.view.df = None;
1360 }
1361
1362 pub(super) fn slice_buffer_into_display(&mut self) {
1364 let full_df = match self.view.buffered_df.as_ref() {
1365 Some(df) => df,
1366 None => return,
1367 };
1368
1369 if self.view.locked_columns_count > 0 {
1370 let locked_names: Vec<&str> = self
1371 .view
1372 .column_order
1373 .iter()
1374 .take(self.view.locked_columns_count)
1375 .map(|s| s.as_str())
1376 .collect();
1377 if let Ok(locked_df) = full_df.select(locked_names) {
1378 self.view.locked_df = Some(locked_df);
1379 }
1380 } else {
1381 self.view.locked_df = None;
1382 }
1383
1384 let scroll_names: Vec<&str> = self
1385 .view
1386 .column_order
1387 .iter()
1388 .skip(self.frozen_shown() + self.termcol_index)
1389 .map(|s| s.as_str())
1390 .collect();
1391 if scroll_names.is_empty() {
1392 self.view.df = None;
1393 } else {
1394 if let Ok(scroll_df) = full_df.select(scroll_names) {
1395 self.view.df = Some(scroll_df);
1396 }
1397 }
1398 }
1399
1400 pub fn wants_to_load_ahead(&self) -> bool {
1403 if self.visible_rows == 0
1404 || self.view.buffered_df.is_none()
1405 || !self.page_on_hand(self.view.start_row)
1406 {
1407 return false;
1408 }
1409 let near = self.proximity();
1410 let view_end = self.view.start_row
1411 + self
1412 .visible_rows
1413 .min(self.num_rows_bound().saturating_sub(self.view.start_row));
1414 let behind = self.view.start_row - self.view.buffered_start_row <= near
1415 && self.view.buffered_start_row > 0;
1416 let ahead = self.view.buffered_end_row - view_end <= near
1417 && self.view.buffered_end_row < self.num_rows_bound();
1418 behind || ahead
1419 }
1420
1421 fn proximity(&self) -> usize {
1424 (self.reach_rows(self.pages_lookahead) / 2).max(self.visible_rows)
1425 }
1426
1427 pub fn buffer_position(&self) -> (u64, usize, usize, usize) {
1429 (
1430 self.len_generation(),
1431 self.view.start_row,
1432 self.view.buffered_start_row,
1433 self.view.buffered_end_row,
1434 )
1435 }
1436
1437 pub(crate) fn page_on_hand(&self, start: usize) -> bool {
1439 let bound = self.num_rows_bound();
1440 let end = start + self.visible_rows.min(bound.saturating_sub(start));
1441 self.view.buffered_df.is_some()
1442 && self.view.buffered_end_row > 0
1443 && start >= self.view.buffered_start_row
1444 && end <= self.view.buffered_end_row
1445 }
1446
1447 pub(crate) fn start_to_draw(&mut self) -> usize {
1450 if self.page_on_hand(self.view.start_row) {
1451 self.view.drawn_start = self.view.start_row;
1452 self.view.start_row
1453 } else if self.page_on_hand(self.view.drawn_start) {
1454 self.view.drawn_start
1455 } else {
1456 self.view.start_row
1457 }
1458 }
1459}