1use std::path::Path;
17
18pub(crate) const READER: crate::formats::readers::Reader = crate::formats::readers::Reader {
20 scan,
21 python: Some(crate::export::python_script::Python {
22 call: "pl.read_database",
23 eager: true,
24 glob_flag: false,
25 arguments: Some(crate::export::python_script::sqlite_arguments),
26 }),
27 signatures: &[crate::formats::readers::Signature {
28 says: |head, _| looks_like(head),
29 kind: crate::formats::readers::Kind::Magic,
30 trusted: crate::formats::readers::Trusted {
31 tables: true,
32 ..crate::formats::readers::EVERYWHERE
33 },
34 }],
35 tables: Some(tables),
36 table_schema: Some(
37 |file, name| match pick(tables(file).ok()?, name, file).ok()? {
38 Pick::One(table) => schema_preview(file, &table),
39 Pick::Several(_) => None,
40 },
41 ),
42 bytes_decide: true,
43 ..crate::formats::readers::BASE
44};
45
46pub const MAGIC: &[u8; 16] = b"SQLite format 3\0";
48
49pub fn looks_like(head: &[u8]) -> bool {
51 head.starts_with(MAGIC)
52}
53
54#[cfg(feature = "sqlite")]
57fn is_internal(name: &str, kind: &str) -> bool {
58 name.get(..7)
59 .is_some_and(|prefix| prefix.eq_ignore_ascii_case("sqlite_"))
60 || kind == "shadow"
61}
62
63#[derive(Debug, Clone, PartialEq, Eq)]
65pub struct Table {
66 pub name: String,
67 pub kind: String,
69 pub internal: bool,
71 pub columns: Vec<(String, String)>,
73}
74
75impl Table {
76 pub fn plain(
78 name: impl Into<String>,
79 kind: &str,
80 columns: impl IntoIterator<Item = impl Into<String>>,
81 ) -> Self {
82 Self {
83 name: name.into(),
84 kind: kind.to_string(),
85 internal: false,
86 columns: columns
87 .into_iter()
88 .map(|c| (c.into(), String::new()))
89 .collect(),
90 }
91 }
92}
93
94#[derive(Debug, Clone, Copy, PartialEq, Eq)]
97pub enum Affinity {
98 Integer,
99 Text,
100 Blob,
102 Real,
103 Numeric,
104}
105
106impl Affinity {
107 pub fn of(declared: &str) -> Self {
108 let upper = declared.to_ascii_uppercase();
109 if upper.contains("INT") {
110 Self::Integer
111 } else if ["CHAR", "CLOB", "TEXT"].iter().any(|t| upper.contains(t)) {
112 Self::Text
113 } else if upper.contains("BLOB") || upper.trim().is_empty() {
114 Self::Blob
115 } else if ["REAL", "FLOA", "DOUB"].iter().any(|t| upper.contains(t)) {
116 Self::Real
117 } else {
118 Self::Numeric
119 }
120 }
121}
122
123#[derive(Debug, Clone, PartialEq, Eq)]
126pub enum Pick {
127 One(Table),
128 Several(Vec<Table>),
129}
130
131pub fn pick(tables: Vec<Table>, wanted: Option<&str>, display: &Path) -> color_eyre::Result<Pick> {
134 crate::formats::members::pick(
135 tables,
136 wanted,
137 display,
138 " --table sqlite_master shows its schema.",
139 )
140}
141
142#[cfg(feature = "sqlite")]
143pub use read::*;
144
145#[cfg(not(feature = "sqlite"))]
146pub use unsupported::*;
147
148pub struct Opened {
150 pub lf: polars::prelude::LazyFrame,
152 pub pushdown: std::sync::Arc<dyn crate::formats::pushdown::Pushdown>,
154 pub hold: Hold,
156 pub other_tables: Vec<String>,
158}
159
160pub struct Hold(std::sync::Arc<std::sync::atomic::AtomicBool>);
162
163impl Drop for Hold {
164 fn drop(&mut self) {
165 self.0.store(true, std::sync::atomic::Ordering::Relaxed);
166 }
167}
168
169#[cfg(not(feature = "sqlite"))]
170mod unsupported {
171 use std::path::Path;
172
173 use color_eyre::Result;
174 use color_eyre::eyre::eyre;
175
176 use super::{Opened, Table};
177
178 fn refused() -> color_eyre::Report {
179 eyre!(
180 "This build of datui reads no SQLite databases: it was built without the sqlite feature."
181 )
182 }
183
184 pub fn tables(_path: &Path) -> Result<Vec<Table>> {
185 Err(refused())
186 }
187
188 pub fn detail(_path: &Path, _tables: &[Table]) -> Option<crate::formats::text_formats::Detail> {
189 None
190 }
191
192 pub fn schema_preview(
193 _path: &Path,
194 _table: &Table,
195 ) -> Option<crate::home::discover::SchemaPreview> {
196 None
197 }
198
199 pub fn open_table(
200 _file: &Path,
201 _display: &Path,
202 _table: &Table,
203 _others: &[Table],
204 ) -> Result<Opened> {
205 Err(refused())
206 }
207}
208
209#[cfg(feature = "sqlite")]
210mod read {
211 use std::collections::{BTreeMap, HashMap};
212 use std::path::{Path, PathBuf};
213 use std::sync::atomic::{AtomicBool, Ordering};
214 use std::sync::{Arc, Mutex, OnceLock};
215 use std::time::{Duration, Instant};
216
217 use color_eyre::Result;
218 use polars::prelude::*;
219 use rusqlite::config::DbConfig;
220 use rusqlite::types::{Value, ValueRef};
221 use rusqlite::{Connection, OpenFlags};
222
223 use super::{Affinity, Hold, Opened, Table, is_internal};
224 use crate::app::modals::filter_modal::{FilterOperator, FilterStatement, LogicalOperator};
225 use crate::error_display::{FileError, file_message};
226 use crate::formats::pushdown::{Counter, Pushdown, PushedView, Windowed};
227 use crate::notes::Note;
228 use crate::numfmt::group_chrome;
229
230 const BUSY: Duration = Duration::from_secs(2);
232 const STEPS_PER_LOOK: i32 = 100_000;
235 const PREVIEW_TIME: Duration = Duration::from_secs(2);
237 const MAX_DESCRIBED: usize = 1000;
239
240 fn open(path: &Path) -> Result<Connection> {
246 let not_a_database = |e: rusqlite::Error| match e.sqlite_error_code() {
247 Some(rusqlite::ErrorCode::NotADatabase) => {
248 FileError::new(path, "not a SQLite database").into()
249 }
250 _ => color_eyre::Report::new(e),
251 };
252 if !(is_wal(path) && !beside(path, "-wal").exists()) {
253 match plain(path) {
254 Ok(conn) => return Ok(conn),
255 Err(e) if cannot_open(&e) && beside(path, "-journal").exists() => {
258 return Err(FileError::new(
259 path,
260 "the database was left mid-write by a program that stopped: its -journal has to be rolled back first, which datui does not do. Opening it once with the sqlite3 tool rolls it back.",
261 )
262 .into());
263 }
264 Err(e) if cannot_open(&e) && beside(path, "-wal").exists() => {
267 return Err(FileError::new(
268 path,
269 "the database has a -wal file that cannot be read from here without a -shm file beside it, and its directory is read only. Copy the database and its -wal to a writable directory.",
270 )
271 .into());
272 }
273 Err(e) if cannot_open(&e) => {}
275 Err(e) => return Err(not_a_database(e)),
276 }
277 }
278 Connection::open_with_flags(immutable_uri(path), FLAGS | OpenFlags::SQLITE_OPEN_URI)
279 .and_then(check)
280 .map_err(not_a_database)
281 }
282
283 const FLAGS: OpenFlags =
284 OpenFlags::SQLITE_OPEN_READ_ONLY.union(OpenFlags::SQLITE_OPEN_NO_MUTEX);
285
286 fn plain(path: &Path) -> rusqlite::Result<Connection> {
287 check(Connection::open_with_flags(path, FLAGS)?)
288 }
289
290 fn is_wal(path: &Path) -> bool {
292 use std::io::Read;
293 let mut header = [0u8; 20];
294 std::fs::File::open(path)
295 .and_then(|mut f| f.read_exact(&mut header))
296 .is_ok()
297 && header[18] == 2
298 && header[19] == 2
299 }
300
301 fn beside(path: &Path, suffix: &str) -> std::path::PathBuf {
303 let mut name = path.as_os_str().to_owned();
304 name.push(suffix);
305 name.into()
306 }
307
308 fn cannot_open(e: &rusqlite::Error) -> bool {
309 matches!(
310 e.sqlite_error_code(),
311 Some(rusqlite::ErrorCode::CannotOpen | rusqlite::ErrorCode::ReadOnly)
312 )
313 }
314
315 fn check(conn: Connection) -> rusqlite::Result<Connection> {
318 conn.busy_timeout(BUSY)?;
319 conn.set_db_config(DbConfig::SQLITE_DBCONFIG_DEFENSIVE, true)?;
320 conn.set_db_config(DbConfig::SQLITE_DBCONFIG_TRUSTED_SCHEMA, false)?;
321 conn.set_db_config(DbConfig::SQLITE_DBCONFIG_ENABLE_TRIGGER, false)?;
323 conn.set_db_config(DbConfig::SQLITE_DBCONFIG_ENABLE_FTS3_TOKENIZER, false)?;
324 conn.pragma_update(None, "query_only", true)?;
325 conn.query_row("SELECT count(*) FROM main.sqlite_schema", [], |_| Ok(()))?;
326 Ok(conn)
327 }
328
329 fn immutable_uri(path: &Path) -> String {
331 let absolute = std::path::absolute(path).unwrap_or_else(|_| path.to_path_buf());
332 let text = absolute.to_string_lossy().replace('\\', "/");
333 let mut uri = String::from(if text.starts_with('/') {
336 "file://"
337 } else {
338 "file:///"
339 });
340 for byte in text.bytes() {
341 match byte {
342 b'?' | b'#' | b'%' | 0x80.. => uri.push_str(&format!("%{byte:02X}")),
343 _ => uri.push(char::from(byte)),
344 }
345 }
346 uri.push_str("?immutable=1");
347 uri
348 }
349
350 fn quoted(name: &str) -> String {
352 format!("\"{}\"", name.replace('"', "\"\""))
353 }
354
355 pub fn tables(path: &Path) -> Result<Vec<Table>> {
358 let conn = open(path)?;
359 let kinds: std::collections::HashMap<String, String> = conn
362 .prepare("SELECT name, type FROM pragma_table_list WHERE schema = 'main'")
363 .and_then(|mut stmt| {
364 stmt.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?
365 .collect::<rusqlite::Result<_>>()
366 })
367 .unwrap_or_default();
368 let mut stmt = conn.prepare(
369 "SELECT name, type FROM main.sqlite_schema \
370 WHERE type IN ('table', 'view') ORDER BY rowid",
371 )?;
372 let listed: Vec<(String, String)> = stmt
373 .query_map([], |row| {
374 Ok((
375 text_of(row.get_ref(0)?).unwrap_or_default(),
376 text_of(row.get_ref(1)?).unwrap_or_default(),
377 ))
378 })?
379 .collect::<rusqlite::Result<_>>()?;
380 let mut tables: Vec<Table> = listed
381 .into_iter()
382 .filter(|(name, _)| !name.is_empty())
383 .map(|(name, kind)| {
384 let kind = kinds.get(&name).cloned().unwrap_or(kind);
385 Table {
386 internal: is_internal(&name, &kind),
387 name,
388 kind,
389 columns: Vec::new(),
390 }
391 })
392 .collect();
393 tables.push(Table {
394 name: "sqlite_master".to_string(),
395 kind: "table".to_string(),
396 internal: true,
397 columns: Vec::new(),
398 });
399 for table in tables.iter_mut().take(MAX_DESCRIBED) {
400 table.columns = columns_of(&conn, &table.name);
401 }
402 Ok(tables)
403 }
404
405 pub fn detail(path: &Path, tables: &[Table]) -> Option<crate::formats::text_formats::Detail> {
409 use crate::formats::model_files::MetaValue;
410 use crate::formats::text_formats::count;
411 let conn = open(path).ok()?;
412 let pragma = |name: &str| -> Option<i64> {
413 conn.query_row(&format!("PRAGMA main.{name}"), [], |row| row.get(0))
414 .ok()
415 };
416 let text = |name: &str| -> Option<String> {
417 conn.query_row(&format!("PRAGMA main.{name}"), [], |row| row.get(0))
418 .ok()
419 };
420 let middot = crate::glyphs::get().middot;
421 let page_size = pragma("page_size").unwrap_or(0);
422 let pages = pragma("page_count").unwrap_or(0);
423 let mut lines = vec![format!(
424 "Page size: {} {middot} {}",
425 group_chrome(usize::try_from(page_size).unwrap_or(0)),
426 count(u64::try_from(pages).unwrap_or(0), "page", "pages"),
427 )];
428 let mut versions = format!("Schema version: {}", pragma("schema_version").unwrap_or(0));
429 if let Some(user) = pragma("user_version").filter(|v| *v != 0) {
430 versions.push_str(&format!(" {middot} user version: {user}"));
431 }
432 if let Some(encoding) = text("encoding") {
433 versions.push_str(&format!(" {middot} {encoding}"));
434 }
435 lines.push(versions);
436 let mut analyzed: std::collections::HashMap<String, u64> = Default::default();
439 let has_stats = tables.iter().any(|t| t.name == "sqlite_stat1");
440 if has_stats
441 && let Ok(mut stmt) = conn.prepare("SELECT tbl, stat FROM main.sqlite_stat1")
442 && let Ok(rows) = stmt.query_map([], |row| {
443 Ok((
444 text_of(row.get_ref(0)?).unwrap_or_default(),
445 text_of(row.get_ref(1)?).unwrap_or_default(),
446 ))
447 })
448 {
449 for (table, stat) in rows.flatten() {
450 if let Some(n) = stat.split(' ').next().and_then(|n| n.parse::<u64>().ok()) {
451 let rows = analyzed.entry(table).or_default();
452 *rows = (*rows).max(n);
453 }
454 }
455 }
456 let own: Vec<&Table> = tables.iter().filter(|t| !t.internal).collect();
457 let views = own.iter().filter(|t| t.kind == "view").count();
458 let mut held = count((own.len() - views) as u64, "table", "tables");
459 if views > 0 {
460 held.push_str(&format!(
461 " {middot} {}",
462 count(views as u64, "view", "views")
463 ));
464 }
465 lines.push(held);
466 lines.push(if analyzed.is_empty() {
467 "Rows: not stored; ANALYZE stores them".to_string()
468 } else {
469 "Rows: as ANALYZE last stored them".to_string()
470 });
471 let list = crate::formats::text_formats::capped_list(
472 own.iter().map(|t| {
473 let mut said = vec![t.kind.clone()];
474 if !t.columns.is_empty() {
475 said.push(count(t.columns.len() as u64, "column", "columns"));
476 }
477 if let Some(rows) = analyzed.get(&t.name) {
478 said.push(count(*rows, "row", "rows"));
479 }
480 (t.name.clone(), MetaValue::Text(said.join(", ")))
481 }),
482 own.len(),
483 );
484 Some(crate::formats::text_formats::Detail {
485 tab: crate::formats::text_formats::tab(crate::FileFormat::Sqlite),
486 lines,
487 list_title: "Tables",
488 tables: own.iter().map(|t| t.name.clone()).collect(),
489 list,
490 ..Default::default()
491 })
492 }
493
494 fn columns_of(conn: &Connection, table: &str) -> Vec<(String, String)> {
497 conn.prepare("SELECT name, type FROM pragma_table_info(?1, 'main')")
498 .and_then(|mut stmt| {
499 stmt.query_map([table], |row| {
500 Ok((
501 text_of(row.get_ref(0)?).unwrap_or_default(),
502 text_of(row.get_ref(1)?).unwrap_or_default(),
503 ))
504 })?
505 .collect::<rusqlite::Result<_>>()
506 })
507 .unwrap_or_default()
508 }
509
510 fn text_of(value: ValueRef<'_>) -> Option<String> {
511 match value {
512 ValueRef::Text(bytes) => Some(String::from_utf8_lossy(bytes).into_owned()),
513 ValueRef::Integer(i) => Some(i.to_string()),
514 _ => None,
515 }
516 }
517
518 const BATCH_ROWS: usize = 1 << 16;
520 const BATCH_CELLS: usize = 1 << 22;
523 const BATCH_BYTES: usize = 64 << 20;
525 const SAMPLE_ROWS: usize = 1000;
527 const SAMPLE_TIME: Duration = Duration::from_secs(2);
529 const CHECKPOINTS: usize = 4096;
531 const MEMO_VIEWS: usize = 16;
533 const CENSUS_COLUMNS: usize = 1000;
535
536 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
538 enum Kind {
539 Int,
540 Float,
541 Text,
542 Bytes,
543 }
544
545 impl Kind {
546 fn dtype(self) -> DataType {
547 match self {
548 Self::Int => DataType::Int64,
549 Self::Float => DataType::Float64,
550 Self::Text => DataType::String,
551 Self::Bytes => DataType::Binary,
552 }
553 }
554
555 fn classes(self) -> &'static str {
557 match self {
558 Self::Int => "'integer', 'null'",
559 Self::Float => "'integer', 'real', 'null'",
560 Self::Text => "'text', 'null'",
561 Self::Bytes => "'blob', 'null'",
562 }
563 }
564
565 fn read(self, e: &str) -> String {
569 match self {
570 Self::Int => format!("CASE WHEN typeof({e}) = 'integer' THEN {e} END"),
571 Self::Float => {
572 format!(
573 "CASE WHEN typeof({e}) IN ('integer', 'real') THEN CAST({e} AS REAL) END"
574 )
575 }
576 Self::Text => format!(
577 "CASE typeof({e}) WHEN 'blob' THEN 'X''' || hex({e}) || '''' ELSE CAST({e} AS TEXT) END"
578 ),
579 Self::Bytes => format!("CAST({e} AS BLOB)"),
580 }
581 }
582 }
583
584 #[derive(Debug, Default, Clone, Copy)]
586 struct Seen {
587 int: bool,
588 real: bool,
589 text: bool,
590 blob: bool,
591 }
592
593 impl Seen {
594 fn add(&mut self, value: ValueRef<'_>) {
595 match value {
596 ValueRef::Null => {}
597 ValueRef::Integer(_) => self.int = true,
598 ValueRef::Real(_) => self.real = true,
599 ValueRef::Text(_) => self.text = true,
600 ValueRef::Blob(_) => self.blob = true,
601 }
602 }
603 }
604
605 fn decide(affinity: Affinity, declared: &str, seen: Seen) -> (Kind, bool) {
609 let numbers = seen.int || seen.real;
610 let several = [numbers, seen.text, seen.blob]
611 .iter()
612 .filter(|&&s| s)
613 .count()
614 > 1;
615 match affinity {
616 Affinity::Integer | Affinity::Real if seen.text || seen.blob => (Kind::Text, true),
617 Affinity::Integer if seen.real => (Kind::Float, false),
618 Affinity::Integer => (Kind::Int, false),
619 Affinity::Real => (Kind::Float, false),
620 Affinity::Text => (Kind::Text, seen.blob),
621 Affinity::Blob if !declared.trim().is_empty() => {
622 if numbers || seen.text {
623 (Kind::Text, true)
624 } else {
625 (Kind::Bytes, false)
626 }
627 }
628 _ if several => (Kind::Text, true),
629 _ if seen.real => (Kind::Float, false),
630 _ if seen.int => (Kind::Int, false),
631 _ if seen.blob => (Kind::Bytes, false),
632 _ if seen.text => (Kind::Text, false),
633 Affinity::Numeric => (Kind::Float, false),
635 _ => (Kind::Text, false),
636 }
637 }
638
639 fn described(
642 stmt: &rusqlite::Statement<'_>,
643 table: &Table,
644 ) -> (Vec<String>, Vec<(Affinity, String)>) {
645 let mut seen = std::collections::HashSet::new();
646 let names = stmt
647 .column_names()
648 .iter()
649 .enumerate()
650 .map(|(i, name)| {
651 let base = if name.is_empty() {
652 format!("column_{}", i + 1)
653 } else {
654 name.to_string()
655 };
656 let mut unique = base.clone();
657 let mut n = 1;
658 while !seen.insert(unique.clone()) {
659 unique = format!("{base}_{n}");
660 n += 1;
661 }
662 unique
663 })
664 .collect::<Vec<_>>();
665 let declared: Vec<String> = if table.columns.len() == names.len() {
668 table.columns.iter().map(|(_, t)| t.clone()).collect()
669 } else {
670 vec![String::new(); names.len()]
671 };
672 let affinities = declared
673 .into_iter()
674 .map(|d| (Affinity::of(&d), d))
675 .collect();
676 (names, affinities)
677 }
678
679 fn stoppable(
681 conn: &Connection,
682 stop: Arc<AtomicBool>,
683 deadline: Option<Instant>,
684 ) -> rusqlite::Result<()> {
685 conn.progress_handler(
686 STEPS_PER_LOOK,
687 Some(move || {
688 stop.load(Ordering::Relaxed) || deadline.is_some_and(|d| Instant::now() > d)
689 }),
690 )
691 }
692
693 fn key_of(conn: &Connection, table: &Table) -> Vec<String> {
696 let quoted_table = quoted(&table.name);
697 let shadowed = |alias: &str| {
698 table
699 .columns
700 .iter()
701 .any(|(name, _)| name.eq_ignore_ascii_case(alias))
702 };
703 if let Some(alias) = ["rowid", "_rowid_", "oid"]
704 .into_iter()
705 .find(|alias| !shadowed(alias))
706 && conn
707 .prepare(&format!("SELECT {alias} FROM main.{quoted_table} LIMIT 0"))
708 .is_ok()
709 {
710 return vec![alias.to_string()];
711 }
712 if table.kind == "view" {
713 return Vec::new();
714 }
715 conn.prepare("SELECT name FROM pragma_table_info(?1, 'main') WHERE pk > 0 ORDER BY pk")
716 .and_then(|mut stmt| {
717 stmt.query_map([&table.name], |row| row.get::<_, String>(0))?
718 .collect::<rusqlite::Result<Vec<_>>>()
719 })
720 .map(|names| names.iter().map(|n| quoted(n)).collect())
721 .unwrap_or_default()
722 }
723
724 struct Source {
726 file: PathBuf,
727 display: PathBuf,
728 with: String,
731 keys: usize,
732 columns: Vec<SourceColumn>,
733 schema: SchemaRef,
734 census: OnceLock<Census>,
735 memo: Mutex<HashMap<String, Memo>>,
736 stop: Arc<AtomicBool>,
738 table: Table,
739 }
740
741 struct SourceColumn {
742 name: PlSmallStr,
743 kind: Kind,
744 mixed: bool,
746 numeric_text: bool,
749 }
750
751 struct Census {
754 rows: usize,
755 misfits: Vec<u64>,
756 }
757
758 #[derive(Default)]
760 struct Memo {
761 rows: Option<usize>,
762 places: BTreeMap<usize, Vec<Value>>,
764 }
765
766 #[derive(Debug, Clone)]
768 struct Atom {
769 column: usize,
770 operator: FilterOperator,
771 value: Value,
772 }
773
774 #[derive(Debug)]
777 struct View {
778 key: String,
780 conditions: Vec<(LogicalOperator, Atom)>,
781 sort: Vec<(usize, bool)>,
783 reversed: bool,
785 }
786
787 impl View {
788 fn new(
789 conditions: Vec<(LogicalOperator, Atom)>,
790 sort: Vec<(usize, bool)>,
791 reversed: bool,
792 ) -> Self {
793 Self {
794 key: format!("{conditions:?} {sort:?} {reversed}"),
795 conditions,
796 sort,
797 reversed,
798 }
799 }
800
801 fn whole() -> Self {
802 Self::new(Vec::new(), Vec::new(), false)
803 }
804
805 fn natural(&self) -> bool {
807 self.sort.is_empty()
808 }
809 }
810
811 struct Query<'a> {
813 columns: &'a [usize],
814 also: Option<(String, Vec<Value>)>,
817 backward: bool,
819 after: Option<Vec<Value>>,
821 limit: Option<usize>,
822 offset: usize,
823 with_keys: bool,
825 }
826
827 impl Source {
828 fn open(file: &Path, display: &Path, table: &Table) -> Result<Self> {
829 let conn = open(file)?;
830 let stop = Arc::new(AtomicBool::new(false));
831 stoppable(&conn, stop.clone(), Some(Instant::now() + SAMPLE_TIME))?;
832 let key = key_of(&conn, table);
833 let star = format!("SELECT * FROM main.{}", quoted(&table.name));
834 let stmt = conn.prepare(&star).map_err(|e| named(e.into(), display))?;
835 let (names, affinities) = described(&stmt, table);
836 drop(stmt);
837 let width = names.len();
838 let mut with = String::from("WITH s(");
839 let mut parts: Vec<String> = (0..key.len()).map(|i| format!("k{i}")).collect();
840 parts.extend((0..width).map(|i| format!("c{i}")));
841 with.push_str(&parts.join(", "));
842 with.push_str(") AS (SELECT ");
843 for k in &key {
844 with.push_str(k);
845 with.push_str(", ");
846 }
847 with.push_str(&format!("* FROM main.{})", quoted(&table.name)));
848
849 let mut seen = vec![Seen::default(); width];
851 let list = (0..width)
852 .map(|i| format!("s.c{i}"))
853 .collect::<Vec<_>>()
854 .join(", ");
855 if width > 0 {
856 let sql = format!("{with} SELECT {list} FROM s LIMIT {SAMPLE_ROWS}");
857 let mut stmt = conn.prepare(&sql).map_err(|e| named(e.into(), display))?;
858 let mut rows = stmt.query([]).map_err(|e| named(e.into(), display))?;
859 loop {
862 match rows.next() {
863 Ok(Some(row)) => {
864 for (i, seen) in seen.iter_mut().enumerate() {
865 seen.add(row.get_ref(i)?);
866 }
867 }
868 Ok(None) => break,
869 Err(e)
870 if e.sqlite_error_code()
871 == Some(rusqlite::ErrorCode::OperationInterrupted) =>
872 {
873 break;
874 }
875 Err(e) => return Err(named(e.into(), display)),
876 }
877 }
878 }
879 let columns: Vec<SourceColumn> = names
880 .into_iter()
881 .zip(affinities)
882 .zip(seen)
883 .map(|((name, (affinity, declared)), seen)| {
884 let (kind, mixed) = decide(affinity, &declared, seen);
885 SourceColumn {
886 name: name.into(),
887 kind,
888 mixed,
889 numeric_text: kind == Kind::Text
890 && matches!(
891 affinity,
892 Affinity::Integer | Affinity::Real | Affinity::Numeric
893 ),
894 }
895 })
896 .collect();
897 let schema = Arc::new(Schema::from_iter(
898 columns
899 .iter()
900 .map(|c| Field::new(c.name.clone(), c.kind.dtype())),
901 ));
902 Ok(Self {
903 file: file.to_path_buf(),
904 display: display.to_path_buf(),
905 with,
906 keys: key.len(),
907 columns,
908 schema,
909 census: OnceLock::new(),
910 memo: Mutex::new(HashMap::new()),
911 stop,
912 table: table.clone(),
913 })
914 }
915
916 fn connect(&self) -> PolarsResult<Connection> {
918 let conn = open(&self.file).map_err(|e| {
919 polars_err!(ComputeError: "{}", crate::error_display::user_message_from_report(&e, Some(&self.display)))
920 })?;
921 stoppable(&conn, self.stop.clone(), None).map_err(|e| self.failed(e))?;
922 Ok(conn)
923 }
924
925 fn failed(&self, e: rusqlite::Error) -> PolarsError {
926 if self.stop.load(Ordering::Relaxed) {
927 polars_err!(ComputeError: "{}", file_message(&self.display, "reading was stopped"))
928 } else {
929 polars_err!(ComputeError: "{}", file_message(&self.display, &e.to_string()))
930 }
931 }
932
933 fn clean(&self, i: usize) -> bool {
936 self.census.get().is_some_and(|c| c.misfits[i] == 0)
937 }
938
939 fn column_sql(&self, i: usize) -> String {
941 let plain = format!("s.c{i}");
942 if self.clean(i) {
943 plain
944 } else {
945 self.columns[i].kind.read(&plain)
946 }
947 }
948
949 fn compared_sql(&self, i: usize) -> String {
952 let e = self.column_sql(i);
953 if self.columns[i].numeric_text {
954 format!("+{e}")
955 } else {
956 e
957 }
958 }
959
960 fn index_of(&self, name: &str) -> Option<usize> {
961 self.columns.iter().position(|c| c.name == name)
962 }
963
964 fn atom(&self, filter: &FilterStatement) -> Option<Atom> {
967 if !filter.operator.takes_value() || filter.operator.is_find() {
970 return None;
971 }
972 let column = self.index_of(&filter.column)?;
973 let kind = self.columns[column].kind;
974 let contains = matches!(
975 filter.operator,
976 FilterOperator::Contains | FilterOperator::NotContains
977 );
978 let value = match kind {
980 Kind::Int if !contains => Value::Integer(filter.value.parse().ok()?),
981 Kind::Float if !contains => Value::Real(filter.value.parse().ok()?),
982 Kind::Text => Value::Text(filter.value.clone()),
983 _ => return None,
984 };
985 Some(Atom {
986 column,
987 operator: filter.operator,
988 value,
989 })
990 }
991
992 fn atom_sql(&self, atom: &Atom, params: &mut Vec<Value>) -> String {
993 let e = self.compared_sql(atom.column);
994 let collate = if self.columns[atom.column].kind == Kind::Text {
996 " COLLATE BINARY"
997 } else {
998 ""
999 };
1000 params.push(atom.value.clone());
1001 match atom.operator {
1002 FilterOperator::Eq => format!("{e} = ?{collate}"),
1003 FilterOperator::NotEq => format!("{e} <> ?{collate}"),
1004 FilterOperator::Gt => format!("{e} > ?{collate}"),
1005 FilterOperator::Lt => format!("{e} < ?{collate}"),
1006 FilterOperator::GtEq => format!("{e} >= ?{collate}"),
1007 FilterOperator::LtEq => format!("{e} <= ?{collate}"),
1008 FilterOperator::Contains => format!("instr({e}, ?) > 0"),
1009 FilterOperator::NotContains => format!("instr({e}, ?) = 0"),
1010 FilterOperator::IsNull | FilterOperator::IsNotNull => {
1012 unreachable!("a null test is not pushed down")
1013 }
1014 FilterOperator::Has | FilterOperator::HasRegex | FilterOperator::HasFuzzy => {
1015 unreachable!("a kept find is not pushed down")
1016 }
1017 }
1018 }
1019
1020 fn conditions_sql(&self, view: &View, params: &mut Vec<Value>) -> Option<String> {
1022 let mut joined: Option<String> = None;
1023 for (logical, atom) in &view.conditions {
1024 let sql = self.atom_sql(atom, params);
1025 joined = Some(match joined {
1026 None => sql,
1027 Some(before) => match logical {
1028 LogicalOperator::And => format!("({before} AND {sql})"),
1029 LogicalOperator::Or => format!("({before} OR {sql})"),
1030 },
1031 });
1032 }
1033 joined
1034 }
1035
1036 fn statement(&self, view: &View, query: &Query<'_>) -> (String, Vec<Value>) {
1037 let mut params = Vec::new();
1038 let mut select: Vec<String> =
1039 query.columns.iter().map(|&i| self.column_sql(i)).collect();
1040 if query.with_keys {
1041 select.extend((0..self.keys).map(|k| format!("s.k{k}")));
1042 }
1043 if select.is_empty() {
1044 select.push("NULL".to_string());
1045 }
1046 let mut sql = format!("{} SELECT {} FROM s", self.with, select.join(", "));
1047 let mut conditions: Vec<String> = Vec::new();
1048 if let Some(c) = self.conditions_sql(view, &mut params) {
1049 conditions.push(c);
1050 }
1051 if let Some((also, also_params)) = &query.also {
1052 conditions.push(also.clone());
1053 params.extend(also_params.iter().cloned());
1054 }
1055 if let Some(after) = &query.after {
1056 let keys: Vec<String> = (0..self.keys).map(|k| format!("s.k{k}")).collect();
1057 params.extend(after.iter().cloned());
1058 let marks = vec!["?"; after.len()];
1059 let past = if view.reversed { "<" } else { ">" };
1060 conditions.push(format!(
1061 "({}) {past} ({})",
1062 keys.join(", "),
1063 marks.join(", ")
1064 ));
1065 }
1066 if !conditions.is_empty() {
1067 sql.push_str(" WHERE ");
1068 sql.push_str(&conditions.join(" AND "));
1069 }
1070 let mut order: Vec<String> = view
1071 .sort
1072 .iter()
1073 .map(|&(i, descending)| {
1074 let collate = if self.columns[i].kind == Kind::Text {
1075 " COLLATE BINARY"
1076 } else {
1077 ""
1078 };
1079 let (direction, nulls) = match descending != query.backward {
1081 true => ("DESC", if query.backward { "FIRST" } else { "LAST" }),
1082 false => ("ASC", if query.backward { "FIRST" } else { "LAST" }),
1083 };
1084 format!("{}{collate} {direction} NULLS {nulls}", self.column_sql(i))
1085 })
1086 .collect();
1087 let keys_backward = view.reversed != query.backward;
1089 order.extend(
1090 (0..self.keys)
1091 .map(|k| format!("s.k{k} {}", if keys_backward { "DESC" } else { "ASC" })),
1092 );
1093 if !order.is_empty() {
1094 sql.push_str(" ORDER BY ");
1095 sql.push_str(&order.join(", "));
1096 }
1097 let limit = query.limit.map_or(-1, |n| n as i64);
1098 sql.push_str(&format!(" LIMIT {limit} OFFSET {}", query.offset));
1099 (sql, params)
1100 }
1101
1102 fn run(
1105 &self,
1106 view: &View,
1107 query: &Query<'_>,
1108 mut each: impl FnMut(DataFrame) -> PolarsResult<bool>,
1109 ) -> PolarsResult<Option<Vec<Value>>> {
1110 let (sql, params) = self.statement(view, query);
1111 let conn = self.connect()?;
1112 let mut stmt = conn.prepare(&sql).map_err(|e| self.failed(e))?;
1113 let mut rows = stmt
1114 .query(rusqlite::params_from_iter(params.iter()))
1115 .map_err(|e| self.failed(e))?;
1116 let width = query.columns.len().max(1);
1117 let batch_rows = (BATCH_CELLS / width).clamp(1, BATCH_ROWS);
1118 let mut batch = Batch::new(self, query.columns);
1119 let mut last_key = None;
1120 loop {
1121 let row = match rows.next() {
1122 Ok(Some(row)) => row,
1123 Ok(None) => break,
1124 Err(e) => return Err(self.failed(e)),
1125 };
1126 for (i, column) in batch.columns.iter_mut().enumerate() {
1127 batch.bytes += column.push(row.get_ref(i).map_err(|e| self.failed(e))?);
1128 }
1129 batch.rows += 1;
1130 if query.with_keys {
1131 let first = query.columns.len();
1132 last_key = Some(
1133 (first..first + self.keys)
1134 .map(|k| row.get::<_, Value>(k))
1135 .collect::<rusqlite::Result<Vec<_>>>()
1136 .map_err(|e| self.failed(e))?,
1137 );
1138 }
1139 if (batch.rows >= batch_rows || batch.bytes >= BATCH_BYTES) && !each(batch.take()?)?
1140 {
1141 return Ok(last_key);
1142 }
1143 }
1144 if batch.rows > 0 || query.columns.is_empty() {
1145 each(batch.take()?)?;
1146 }
1147 Ok(last_key)
1148 }
1149
1150 fn known_rows(&self, view: &View) -> Option<usize> {
1152 if view.conditions.is_empty()
1153 && let Some(census) = self.census.get()
1154 {
1155 return Some(census.rows);
1156 }
1157 self.memo.lock().ok()?.get(&view.key)?.rows
1158 }
1159
1160 fn remember(&self, view: &View, learn: impl FnOnce(&mut Memo)) {
1161 let Ok(mut memo) = self.memo.lock() else {
1162 return;
1163 };
1164 if !memo.contains_key(&view.key) && memo.len() >= MEMO_VIEWS {
1165 memo.clear();
1166 }
1167 learn(memo.entry(view.key.clone()).or_default());
1168 }
1169
1170 fn count(&self, view: &View) -> PolarsResult<usize> {
1173 if let Some(rows) = self.known_rows(view) {
1174 return Ok(rows);
1175 }
1176 let mut params = Vec::new();
1177 let mut sql = format!("{} SELECT count(*) FROM s", self.with);
1178 if let Some(c) = self.conditions_sql(view, &mut params) {
1179 sql.push_str(" WHERE ");
1180 sql.push_str(&c);
1181 }
1182 let conn = self.connect()?;
1183 let rows: i64 = conn
1184 .query_row(&sql, rusqlite::params_from_iter(params.iter()), |row| {
1185 row.get(0)
1186 })
1187 .map_err(|e| self.failed(e))?;
1188 let rows = usize::try_from(rows).unwrap_or(0);
1189 self.remember(view, |memo| memo.rows = Some(rows));
1190 Ok(rows)
1191 }
1192
1193 fn window(
1197 &self,
1198 view: &View,
1199 start: usize,
1200 len: usize,
1201 columns: &[usize],
1202 ) -> PolarsResult<DataFrame> {
1203 let empty = || self.empty(columns);
1204 let total = self.known_rows(view);
1205 if len == 0 || total.is_some_and(|t| start >= t) {
1206 return Ok(empty());
1207 }
1208 let mut frames = Vec::new();
1209 if self.keys > 0
1210 && let Some(total) = total
1211 && start > total / 2
1212 {
1213 let len = len.min(total - start);
1214 let query = Query {
1215 columns,
1216 also: None,
1217 backward: true,
1218 after: None,
1219 limit: Some(len),
1220 offset: total - start - len,
1221 with_keys: false,
1222 };
1223 self.run(view, &query, |df| {
1224 frames.push(df);
1225 Ok(true)
1226 })?;
1227 let df = concat_frames(frames, empty())?;
1228 return Ok(df.reverse());
1229 }
1230 let mut offset = start;
1231 let mut after = None;
1232 let paged = self.keys > 0 && view.natural();
1233 if paged
1234 && let Ok(memo) = self.memo.lock()
1235 && let Some((place, key)) = memo
1236 .get(&view.key)
1237 .and_then(|m| m.places.range(..=start).next_back())
1238 {
1239 offset = start - place;
1240 after = Some(key.clone());
1241 }
1242 let query = Query {
1243 columns,
1244 also: None,
1245 backward: false,
1246 after,
1247 limit: Some(len),
1248 offset,
1249 with_keys: paged,
1250 };
1251 let last = self.run(view, &query, |df| {
1252 frames.push(df);
1253 Ok(true)
1254 })?;
1255 let df = concat_frames(frames, empty())?;
1256 let read = df.height();
1257 if let Some(key) = last {
1258 self.remember(view, |memo| {
1259 if memo.places.len() >= CHECKPOINTS {
1260 memo.places.clear();
1261 }
1262 memo.places.insert(start + read, key);
1263 });
1264 }
1265 if read < len && (start == 0 || read > 0) {
1267 self.remember(view, |memo| memo.rows = Some(start + read));
1268 }
1269 Ok(df)
1270 }
1271
1272 fn whole(
1275 &self,
1276 view: &View,
1277 columns: &[usize],
1278 predicate: Option<&Expr>,
1279 n_rows: Option<usize>,
1280 ) -> PolarsResult<DataFrame> {
1281 let mut read: Vec<usize> = columns.to_vec();
1282 let also = predicate.and_then(|p| {
1283 let mut params = Vec::new();
1284 self.predicate_sql(p, &mut params).map(|sql| (sql, params))
1285 });
1286 if let Some(predicate) = predicate {
1287 for name in predicate.clone().meta().root_names() {
1288 if let Some(i) = self.index_of(&name)
1289 && !read.contains(&i)
1290 {
1291 read.push(i);
1292 }
1293 }
1294 }
1295 let names: Vec<PlSmallStr> = columns
1296 .iter()
1297 .map(|&i| self.columns[i].name.clone())
1298 .collect();
1299 let mut frames = Vec::new();
1300 let mut kept = 0usize;
1301 let wanted = n_rows.unwrap_or(usize::MAX);
1302 let query = Query {
1303 columns: &read,
1304 also,
1305 backward: false,
1306 after: None,
1307 limit: predicate.is_none().then_some(n_rows).flatten(),
1308 offset: 0,
1309 with_keys: false,
1310 };
1311 self.run(view, &query, |df| {
1312 let df = match predicate {
1313 Some(p) => df.lazy().filter(p.clone()).collect()?,
1314 None => df,
1315 };
1316 let df = df.select(names.iter().cloned())?;
1317 kept += df.height();
1318 frames.push(df);
1319 Ok(kept < wanted)
1320 })?;
1321 let mut df = concat_frames(frames, self.empty(columns))?;
1322 if df.height() > wanted {
1323 df = df.head(Some(wanted));
1324 }
1325 Ok(df)
1326 }
1327
1328 fn empty(&self, columns: &[usize]) -> DataFrame {
1330 DataFrame::empty_with_schema(&Schema::from_iter(
1331 columns.iter().map(|&i| {
1332 Field::new(self.columns[i].name.clone(), self.columns[i].kind.dtype())
1333 }),
1334 ))
1335 }
1336 }
1337
1338 impl Source {
1339 fn predicate_sql(&self, e: &Expr, params: &mut Vec<Value>) -> Option<String> {
1342 let Expr::BinaryExpr { left, op, right } = e else {
1343 return None;
1344 };
1345 match op {
1346 Operator::And | Operator::LogicalAnd => {
1347 let mut left_params = Vec::new();
1348 let mut right_params = Vec::new();
1349 let l = self.predicate_sql(left, &mut left_params);
1350 let r = self.predicate_sql(right, &mut right_params);
1351 match (l, r) {
1353 (Some(l), Some(r)) => {
1354 params.extend(left_params);
1355 params.extend(right_params);
1356 Some(format!("({l} AND {r})"))
1357 }
1358 (Some(l), None) => {
1359 params.extend(left_params);
1360 Some(l)
1361 }
1362 (None, Some(r)) => {
1363 params.extend(right_params);
1364 Some(r)
1365 }
1366 (None, None) => None,
1367 }
1368 }
1369 Operator::Or | Operator::LogicalOr => {
1370 let mut both = Vec::new();
1371 let l = self.predicate_sql(left, &mut both)?;
1372 let r = self.predicate_sql(right, &mut both)?;
1373 params.extend(both);
1374 Some(format!("({l} OR {r})"))
1375 }
1376 Operator::Eq
1377 | Operator::NotEq
1378 | Operator::Lt
1379 | Operator::LtEq
1380 | Operator::Gt
1381 | Operator::GtEq => {
1382 let (name, value, op) = match (&**left, &**right) {
1384 (Expr::Column(name), Expr::Literal(value)) => (name, value, *op),
1385 (Expr::Literal(value), Expr::Column(name)) => (name, value, flipped(*op)),
1386 _ => return None,
1387 };
1388 let i = self.index_of(name)?;
1389 let kind = self.columns[i].kind;
1390 let value = literal(value, kind)?;
1391 let sign = match op {
1392 Operator::Eq => "=",
1393 Operator::NotEq => "<>",
1394 Operator::Lt => "<",
1395 Operator::LtEq => "<=",
1396 Operator::Gt => ">",
1397 _ => ">=",
1398 };
1399 let collate = if kind == Kind::Text {
1400 " COLLATE BINARY"
1401 } else {
1402 ""
1403 };
1404 params.push(value);
1405 Some(format!("{} {sign} ?{collate}", self.compared_sql(i)))
1406 }
1407 _ => None,
1408 }
1409 }
1410
1411 fn take_census(&self) -> PolarsResult<()> {
1415 let conn = self.connect()?;
1416 let mut misfits = vec![0u64; self.columns.len()];
1417 let mut rows = 0usize;
1418 let chunks: Vec<Vec<usize>> = (0..self.columns.len().max(1))
1419 .collect::<Vec<_>>()
1420 .chunks(CENSUS_COLUMNS)
1421 .map(<[usize]>::to_vec)
1422 .collect();
1423 for chunk in chunks {
1424 let mut select = vec!["count(*)".to_string()];
1425 select.extend(chunk.iter().filter(|&&i| i < self.columns.len()).map(|&i| {
1426 format!(
1427 "total(typeof(s.c{i}) NOT IN ({}))",
1428 self.columns[i].kind.classes()
1429 )
1430 }));
1431 let sql = format!("{} SELECT {} FROM s", self.with, select.join(", "));
1432 conn.query_row(&sql, [], |row| {
1433 rows = usize::try_from(row.get::<_, i64>(0)?).unwrap_or(0);
1434 for (n, &i) in chunk.iter().enumerate() {
1435 if i < misfits.len() {
1436 misfits[i] = row.get::<_, f64>(n + 1)? as u64;
1437 }
1438 }
1439 Ok(())
1440 })
1441 .map_err(|e| self.failed(e))?;
1442 }
1443 let _ = self.census.set(Census { rows, misfits });
1444 Ok(())
1445 }
1446
1447 fn notes(&self) -> Vec<Note> {
1450 let census = self.census.get();
1451 let misfit = |i: usize| census.map_or(0, |c| c.misfits[i]);
1452 let scope = format!("of the {} {}", self.table.kind, self.table.name);
1453 let note = |summary: String| Note {
1454 summary,
1455 scope: scope.clone(),
1456 read_as_text: None,
1457 passed_over: None,
1458 };
1459 let mut notes = Vec::new();
1460 let mixed: Vec<&str> = self
1461 .columns
1462 .iter()
1463 .enumerate()
1464 .filter(|(i, c)| c.kind == Kind::Text && (c.mixed || misfit(*i) > 0))
1465 .map(|(_, c)| c.name.as_str())
1466 .collect();
1467 if !mixed.is_empty() {
1468 notes.push(note(format!(
1469 "{}: mixed types, read as text",
1470 mixed.join(", ")
1471 )));
1472 }
1473 for (i, column) in self.columns.iter().enumerate() {
1474 let n = misfit(i);
1475 if n == 0 {
1476 continue;
1477 }
1478 let values = format!(
1479 "{} {}",
1480 group_chrome(n as usize),
1481 if n == 1 { "value" } else { "values" }
1482 );
1483 match column.kind {
1484 Kind::Int => notes.push(note(format!(
1485 "{}: {values} not whole numbers, read as null",
1486 column.name
1487 ))),
1488 Kind::Float => notes.push(note(format!(
1489 "{}: {values} not numbers, read as null",
1490 column.name
1491 ))),
1492 Kind::Bytes => notes.push(note(format!(
1493 "{}: {values} not blobs, read as their bytes",
1494 column.name
1495 ))),
1496 Kind::Text => {}
1497 }
1498 }
1499 notes
1500 }
1501 }
1502
1503 fn flipped(op: Operator) -> Operator {
1504 match op {
1505 Operator::Lt => Operator::Gt,
1506 Operator::LtEq => Operator::GtEq,
1507 Operator::Gt => Operator::Lt,
1508 Operator::GtEq => Operator::LtEq,
1509 other => other,
1510 }
1511 }
1512
1513 fn literal(value: &LiteralValue, kind: Kind) -> Option<Value> {
1515 if !value.is_scalar() {
1517 return None;
1518 }
1519 let any = value.to_any_value()?.into_static();
1520 match (kind, any) {
1521 (Kind::Int | Kind::Float, AnyValue::Int8(v)) => Some(Value::Integer(v.into())),
1522 (Kind::Int | Kind::Float, AnyValue::Int16(v)) => Some(Value::Integer(v.into())),
1523 (Kind::Int | Kind::Float, AnyValue::Int32(v)) => Some(Value::Integer(v.into())),
1524 (Kind::Int | Kind::Float, AnyValue::Int64(v)) => Some(Value::Integer(v)),
1525 (Kind::Int | Kind::Float, AnyValue::UInt8(v)) => Some(Value::Integer(v.into())),
1526 (Kind::Int | Kind::Float, AnyValue::UInt16(v)) => Some(Value::Integer(v.into())),
1527 (Kind::Int | Kind::Float, AnyValue::UInt32(v)) => Some(Value::Integer(v.into())),
1528 (Kind::Int | Kind::Float, AnyValue::Float64(v)) if !v.is_nan() => Some(Value::Real(v)),
1529 (Kind::Int | Kind::Float, AnyValue::Float32(v)) if !v.is_nan() => {
1530 Some(Value::Real(v.into()))
1531 }
1532 (Kind::Text, AnyValue::String(s)) => Some(Value::Text(s.to_string())),
1533 (Kind::Text, AnyValue::StringOwned(s)) => Some(Value::Text(s.to_string())),
1534 _ => None,
1535 }
1536 }
1537
1538 fn concat_frames(frames: Vec<DataFrame>, empty: DataFrame) -> PolarsResult<DataFrame> {
1540 let mut frames = frames.into_iter();
1541 let Some(mut df) = frames.next() else {
1542 return Ok(empty);
1543 };
1544 for next in frames {
1545 df.vstack_mut_owned(next)?;
1546 }
1547 df.rechunk_mut();
1548 Ok(df)
1549 }
1550
1551 struct Batch {
1553 columns: Vec<Buffer>,
1554 names: Vec<PlSmallStr>,
1555 rows: usize,
1556 bytes: usize,
1557 }
1558
1559 enum Buffer {
1560 Int(Vec<Option<i64>>),
1561 Float(Vec<Option<f64>>),
1562 Text(Vec<Option<String>>),
1563 Bytes(Vec<Option<Vec<u8>>>),
1564 }
1565
1566 impl Batch {
1567 fn new(source: &Source, columns: &[usize]) -> Self {
1568 Self {
1569 columns: columns
1570 .iter()
1571 .map(|&i| match source.columns[i].kind {
1572 Kind::Int => Buffer::Int(Vec::new()),
1573 Kind::Float => Buffer::Float(Vec::new()),
1574 Kind::Text => Buffer::Text(Vec::new()),
1575 Kind::Bytes => Buffer::Bytes(Vec::new()),
1576 })
1577 .collect(),
1578 names: columns
1579 .iter()
1580 .map(|&i| source.columns[i].name.clone())
1581 .collect(),
1582 rows: 0,
1583 bytes: 0,
1584 }
1585 }
1586
1587 fn take(&mut self) -> PolarsResult<DataFrame> {
1588 let height = self.rows;
1589 let columns: Vec<Column> = self
1590 .names
1591 .iter()
1592 .zip(self.columns.iter_mut())
1593 .map(|(name, buffer)| buffer.take(name.clone()).into())
1594 .collect();
1595 self.rows = 0;
1596 self.bytes = 0;
1597 DataFrame::new(height, columns)
1598 }
1599 }
1600
1601 impl Buffer {
1602 fn push(&mut self, value: ValueRef<'_>) -> usize {
1605 match (self, value) {
1606 (Self::Int(v), ValueRef::Integer(i)) => v.push(Some(i)),
1607 (Self::Int(v), _) => v.push(None),
1608 (Self::Float(v), ValueRef::Integer(i)) => v.push(Some(i as f64)),
1609 (Self::Float(v), ValueRef::Real(f)) => v.push(Some(f)),
1610 (Self::Float(v), _) => v.push(None),
1611 (Self::Text(v), ValueRef::Text(b)) => {
1612 v.push(Some(String::from_utf8_lossy(b).into_owned()));
1613 return b.len();
1614 }
1615 (Self::Text(v), ValueRef::Integer(i)) => v.push(Some(i.to_string())),
1616 (Self::Text(v), ValueRef::Real(f)) => v.push(Some(f.to_string())),
1617 (Self::Text(v), ValueRef::Blob(b)) => {
1618 v.push(Some(blob_text(b)));
1619 return 2 * b.len();
1620 }
1621 (Self::Text(v), ValueRef::Null) => v.push(None),
1622 (Self::Bytes(v), ValueRef::Blob(b) | ValueRef::Text(b)) => {
1623 v.push(Some(b.to_vec()));
1624 return b.len();
1625 }
1626 (Self::Bytes(v), _) => v.push(None),
1627 }
1628 0
1629 }
1630
1631 fn take(&mut self, name: PlSmallStr) -> Series {
1632 match self {
1633 Self::Int(v) => Series::new(name, std::mem::take(v)),
1634 Self::Float(v) => Series::new(name, std::mem::take(v)),
1635 Self::Text(v) => Series::new(name, std::mem::take(v)),
1636 Self::Bytes(v) => {
1637 BinaryChunked::from_iter_options(name, std::mem::take(v).into_iter())
1638 .into_series()
1639 }
1640 }
1641 }
1642 }
1643
1644 fn blob_text(b: &[u8]) -> String {
1646 use std::fmt::Write as _;
1647 let mut text = String::with_capacity(3 + 2 * b.len());
1648 text.push_str("X'");
1649 for byte in b {
1650 let _ = write!(text, "{byte:02X}");
1651 }
1652 text.push('\'');
1653 text
1654 }
1655
1656 pub const SCAN_NAME: &str = "SQLITE";
1658
1659 struct Scan {
1661 source: Arc<Source>,
1662 view: Arc<View>,
1663 window: Option<(usize, usize)>,
1664 }
1665
1666 impl Scan {
1667 fn into_lazy(self) -> PolarsResult<LazyFrame> {
1668 let schema = self.source.schema.clone();
1669 LazyFrame::anonymous_scan(
1670 Arc::new(self),
1671 ScanArgsAnonymous {
1672 schema: Some(schema),
1673 name: SCAN_NAME,
1674 ..Default::default()
1675 },
1676 )
1677 }
1678 }
1679
1680 impl AnonymousScan for Scan {
1681 fn as_any(&self) -> &dyn std::any::Any {
1682 self
1683 }
1684
1685 fn schema(&self, _infer_schema_length: Option<usize>) -> PolarsResult<SchemaRef> {
1686 Ok(self.source.schema.clone())
1687 }
1688
1689 fn allows_projection_pushdown(&self) -> bool {
1690 true
1691 }
1692
1693 fn allows_predicate_pushdown(&self) -> bool {
1696 self.window.is_none()
1697 }
1698
1699 fn scan(&self, mut args: AnonymousScanArgs) -> PolarsResult<DataFrame> {
1700 args.predicate = crate::formats::pushdown::evaluable(args.predicate.take());
1701 let columns: Vec<usize> = match &args.with_columns {
1702 Some(names) => names
1703 .iter()
1704 .map(|name| {
1705 self.source
1706 .index_of(name)
1707 .ok_or_else(|| polars_err!(ColumnNotFound: "{name}"))
1708 })
1709 .collect::<PolarsResult<_>>()?,
1710 None => (0..self.source.columns.len()).collect(),
1711 };
1712 match self.window {
1713 Some((start, len)) => {
1714 let len = args.n_rows.map_or(len, |n| n.min(len));
1715 self.source.window(&self.view, start, len, &columns)
1716 }
1717 None => {
1718 self.source
1719 .whole(&self.view, &columns, args.predicate.as_ref(), args.n_rows)
1720 }
1721 }
1722 }
1723 }
1724
1725 struct Windows {
1727 source: Arc<Source>,
1728 view: Arc<View>,
1729 }
1730
1731 impl Windowed for Windows {
1732 fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
1733 Scan {
1734 source: self.source.clone(),
1735 view: self.view.clone(),
1736 window: Some((start, len)),
1737 }
1738 .into_lazy()
1739 }
1740 }
1741
1742 struct InPlace(Arc<Source>);
1744
1745 impl InPlace {
1746 fn pushed(&self, view: View) -> Option<PushedView> {
1747 let view = Arc::new(view);
1748 let source = self.0.clone();
1749 let lf = Scan {
1750 source: source.clone(),
1751 view: view.clone(),
1752 window: None,
1753 }
1754 .into_lazy()
1755 .ok()?;
1756 let counter: Counter = {
1757 let (source, view) = (source.clone(), view.clone());
1758 Arc::new(move || source.count(&view))
1759 };
1760 Some(PushedView {
1761 lf,
1762 window: Arc::new(Windows { source, view }),
1763 counter,
1764 })
1765 }
1766 }
1767
1768 impl Pushdown for InPlace {
1769 fn view(
1770 &self,
1771 filters: &[FilterStatement],
1772 sort: &[(String, bool)],
1773 reversed: bool,
1774 ) -> Option<PushedView> {
1775 let source = &self.0;
1776 let conditions = filters
1777 .iter()
1778 .map(|f| Some((f.logical_op, source.atom(f)?)))
1779 .collect::<Option<Vec<_>>>()?;
1780 let sort = sort
1781 .iter()
1782 .map(|(name, descending)| Some((source.index_of(name)?, *descending)))
1783 .collect::<Option<Vec<_>>>()?;
1784 let reversed = reversed && sort.is_empty();
1785 if reversed && source.keys == 0 {
1787 return None;
1788 }
1789 self.pushed(View::new(conditions, sort, reversed))
1790 }
1791
1792 fn notes(&self) -> Vec<Note> {
1793 self.0.notes()
1794 }
1795
1796 #[cfg(test)]
1797 fn as_any(&self) -> &dyn std::any::Any {
1798 self
1799 }
1800 }
1801
1802 pub fn open_table(
1806 file: &Path,
1807 display: &Path,
1808 table: &Table,
1809 others: &[Table],
1810 ) -> Result<Opened> {
1811 open_table_with(file, display, table, others, true)
1812 }
1813
1814 pub(super) fn open_table_with(
1817 file: &Path,
1818 display: &Path,
1819 table: &Table,
1820 others: &[Table],
1821 census: bool,
1822 ) -> Result<Opened> {
1823 let source = Arc::new(Source::open(file, display, table)?);
1824 let hold = Hold(source.stop.clone());
1825 if census {
1827 let source = source.clone();
1828 std::thread::Builder::new()
1829 .name("sqlite-census".to_string())
1830 .spawn(move || {
1831 if let Err(e) = source.take_census() {
1832 log::debug!(target: "datui", "sqlite census: {e}");
1833 }
1834 })?;
1835 }
1836 let table_source = InPlace(source);
1837 let whole = table_source
1838 .pushed(View::whole())
1839 .ok_or_else(|| FileError::new(display, format!("could not read \"{}\"", table.name)))?;
1840 Ok(Opened {
1841 lf: whole.lf,
1842 pushdown: Arc::new(table_source),
1843 hold,
1844 other_tables: other_tables(table, others),
1845 })
1846 }
1847
1848 pub fn schema_preview(
1851 path: &Path,
1852 table: &Table,
1853 ) -> Option<crate::home::discover::SchemaPreview> {
1854 const SAMPLE: usize = 100;
1855 let conn = open(path).ok()?;
1856 let stop = Arc::new(AtomicBool::new(false));
1857 stoppable(&conn, stop, Some(Instant::now() + PREVIEW_TIME)).ok()?;
1858 let sql = format!("SELECT * FROM main.{} LIMIT {SAMPLE}", quoted(&table.name));
1859 let mut stmt = conn.prepare(&sql).ok()?;
1860 let (names, affinities) = described(&stmt, table);
1861 let mut seen = vec![Seen::default(); names.len()];
1862 let mut rows = stmt.query([]).ok()?;
1863 while let Some(row) = rows.next().ok()? {
1864 for (i, seen) in seen.iter_mut().enumerate() {
1865 seen.add(row.get_ref(i).ok()?);
1866 }
1867 }
1868 Some(
1869 names
1870 .into_iter()
1871 .zip(affinities)
1872 .zip(seen)
1873 .map(|((name, (affinity, declared)), seen)| {
1874 (name, decide(affinity, &declared, seen).0.dtype())
1875 })
1876 .collect(),
1877 )
1878 }
1879
1880 fn named(e: color_eyre::Report, display: &Path) -> color_eyre::Report {
1882 crate::error_display::in_file(display, e)
1883 }
1884
1885 fn other_tables(table: &Table, others: &[Table]) -> Vec<String> {
1887 const SHOWN: usize = 12;
1888 let own: Vec<&str> = others
1889 .iter()
1890 .filter(|t| !t.internal && t.name != table.name)
1891 .map(|t| t.name.as_str())
1892 .collect();
1893 let mut shown: Vec<String> = own.iter().take(SHOWN).map(|s| s.to_string()).collect();
1894 if own.len() > SHOWN {
1895 shown.push(format!("{} more", group_chrome(own.len() - SHOWN)));
1896 }
1897 shown
1898 }
1899
1900 #[cfg(test)]
1901 pub(super) fn open_for_tests(path: &Path) -> Result<Connection> {
1902 open(path)
1903 }
1904
1905 #[cfg(test)]
1907 fn source_of(opened: &Opened) -> &Arc<Source> {
1908 &opened
1909 .pushdown
1910 .as_any()
1911 .downcast_ref::<InPlace>()
1912 .expect("a table in place")
1913 .0
1914 }
1915
1916 #[cfg(test)]
1918 pub(super) fn census_for_tests(opened: &Opened) {
1919 source_of(opened).take_census().unwrap();
1920 }
1921
1922 #[cfg(test)]
1925 pub(super) fn plan_for_tests(
1926 opened: &Opened,
1927 filters: &[FilterStatement],
1928 sort: &[(String, bool)],
1929 ) -> (String, String) {
1930 let source = source_of(opened);
1931 let conditions = filters
1932 .iter()
1933 .map(|f| (f.logical_op, source.atom(f).unwrap()))
1934 .collect();
1935 let sort = sort
1936 .iter()
1937 .map(|(name, d)| (source.index_of(name).unwrap(), *d))
1938 .collect();
1939 let view = View::new(conditions, sort, false);
1940 let columns: Vec<usize> = (0..source.columns.len()).collect();
1941 let query = Query {
1942 columns: &columns,
1943 also: None,
1944 backward: false,
1945 after: None,
1946 limit: Some(10),
1947 offset: 0,
1948 with_keys: false,
1949 };
1950 let (sql, params) = source.statement(&view, &query);
1951 let conn = source.connect().unwrap();
1952 let mut stmt = conn.prepare(&format!("EXPLAIN QUERY PLAN {sql}")).unwrap();
1953 let plan: Vec<String> = stmt
1954 .query_map(rusqlite::params_from_iter(params.iter()), |row| {
1955 row.get::<_, String>(3)
1956 })
1957 .unwrap()
1958 .collect::<rusqlite::Result<_>>()
1959 .unwrap();
1960 (sql, plan.join("\n"))
1961 }
1962
1963 #[cfg(test)]
1964 pub(super) fn immutable_uri_for_tests(path: &Path) -> String {
1965 immutable_uri(path)
1966 }
1967}
1968
1969#[cfg(all(test, feature = "sqlite"))]
1970mod tests;
1971
1972fn scan(
1975 input: crate::formats::readers::ScanIn<'_>,
1976) -> color_eyre::Result<crate::loading::scan::Scan> {
1977 let file = input.path();
1978 let tables = tables(file)?;
1979 match pick(tables.clone(), input.options.table.as_deref(), file)? {
1980 Pick::One(table) => {
1981 let opened = open_table(file, file, &table, &tables)?;
1982 let detail = detail(file, &tables).map(|detail| {
1983 std::sync::Arc::new(crate::formats::text_formats::Detail {
1984 table: Some(table.name.clone()),
1985 ..detail
1986 })
1987 });
1988 input.report.table = Some(table.name.clone());
1990 input.report.opened = Some(std::sync::Arc::new(crate::formats::members::Opened {
1991 detail,
1992 ..Default::default()
1993 }));
1994 input.report.sqlite = Some(std::sync::Arc::new(crate::SqliteOpen {
1995 pushdown: opened.pushdown,
1996 hold: std::sync::Mutex::new(Some(opened.hold)),
1997 other_tables: opened.other_tables,
1998 }));
1999 Ok(opened.lf.into())
2000 }
2001 Pick::Several(tables) => Ok(crate::formats::members::several(&input, tables)),
2002 }
2003}