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