use crate::array_util::{as_i64, as_str};
use crate::reader::{execute_with_reader, ColumnInfo, Reader, Spec, SqlDialect, TableInfo};
use crate::{naming, DataFrame, Result};
use arrow::array::Array;
use std::cell::{Cell, RefCell};
use std::collections::hash_map::DefaultHasher;
use std::collections::HashSet;
use std::hash::Hasher;
use std::time::{SystemTime, UNIX_EPOCH};
#[derive(Debug, Clone)]
pub struct CacheConfig {
pub enabled: bool,
pub ttl_secs: u64,
pub max_bytes: u64,
}
impl Default for CacheConfig {
fn default() -> Self {
Self {
enabled: true,
ttl_secs: 300,
max_bytes: 512 * 1024 * 1024,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct CacheConfigOverride {
pub enabled: Option<bool>,
pub ttl_secs: Option<u64>,
pub max_bytes: Option<u64>,
}
impl CacheConfig {
pub fn from_env() -> Self {
let mut cfg = Self::default();
if std::env::var("GGSQL_CACHE_DISABLED")
.ok()
.filter(|v| !v.is_empty() && v != "0")
.is_some()
{
cfg.enabled = false;
}
if let Ok(v) = std::env::var("GGSQL_CACHE_TTL") {
if let Ok(secs) = v.trim().parse::<u64>() {
cfg.ttl_secs = secs;
}
}
if let Ok(v) = std::env::var("GGSQL_CACHE_MAX_BYTES") {
if let Some(bytes) = parse_human_bytes(&v) {
cfg.max_bytes = bytes;
}
}
cfg
}
pub fn merge(self, over: CacheConfigOverride) -> Self {
Self {
enabled: over.enabled.unwrap_or(self.enabled),
ttl_secs: over.ttl_secs.unwrap_or(self.ttl_secs),
max_bytes: over.max_bytes.unwrap_or(self.max_bytes),
}
}
}
pub fn parse_human_bytes(s: &str) -> Option<u64> {
let s = s.trim();
let lower = s.to_ascii_lowercase();
let (num, mult) = if let Some(n) = lower.strip_suffix("gb") {
(n, 1024 * 1024 * 1024)
} else if let Some(n) = lower.strip_suffix("mb") {
(n, 1024 * 1024)
} else if let Some(n) = lower.strip_suffix("kb") {
(n, 1024)
} else {
(lower.as_str(), 1)
};
num.trim().parse::<u64>().ok().map(|n| n * mult)
}
fn now_ms() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(i64::MAX)
}
fn estimate_bytes(df: &DataFrame) -> i64 {
df.get_columns()
.iter()
.map(|col| col.get_array_memory_size())
.sum::<usize>() as i64
}
struct MemoEntry {
table_name: String,
fetched_at_epoch_ms: i64,
}
pub struct CachingReader {
primary: Box<dyn Reader + Send>,
cache: Box<dyn Reader + Send>,
primary_uri: String,
cache_scheme: String,
config: CacheConfig,
meta_ready: Cell<bool>,
resident: RefCell<HashSet<String>>,
}
impl CachingReader {
pub fn new(
primary: Box<dyn Reader + Send>,
cache: Box<dyn Reader + Send>,
primary_uri: impl Into<String>,
cache_scheme: impl Into<String>,
) -> Self {
Self::with_config(
primary,
cache,
primary_uri,
cache_scheme,
CacheConfig::from_env(),
)
}
pub fn with_config(
primary: Box<dyn Reader + Send>,
cache: Box<dyn Reader + Send>,
primary_uri: impl Into<String>,
cache_scheme: impl Into<String>,
config: CacheConfig,
) -> Self {
Self {
primary,
cache,
primary_uri: primary_uri.into(),
cache_scheme: cache_scheme.into(),
config,
meta_ready: Cell::new(false),
resident: RefCell::new(HashSet::new()),
}
}
pub fn cache_config(&self) -> &CacheConfig {
&self.config
}
fn cache_key(&self, sql: &str) -> String {
let mut hasher = DefaultHasher::new();
hasher.write(self.primary_uri.as_bytes());
hasher.write(b"\n");
hasher.write(sql.as_bytes());
format!("{:016x}", hasher.finish())
}
fn ensure_meta_table(&self) -> Result<()> {
if self.meta_ready.get() {
return Ok(());
}
let sql = format!(
"CREATE TABLE IF NOT EXISTS {} (\
cache_key VARCHAR PRIMARY KEY, sql VARCHAR NOT NULL, table_name VARCHAR NOT NULL, \
fetched_at_epoch_ms BIGINT NOT NULL, last_accessed_epoch_ms BIGINT NOT NULL, \
byte_estimate BIGINT NOT NULL, row_count BIGINT NOT NULL)",
naming::quote_ident(naming::CACHE_META_TABLE)
);
self.cache.execute_sql(&sql)?;
self.meta_ready.set(true);
Ok(())
}
fn lookup_memo(&self, key: &str) -> Result<Option<MemoEntry>> {
let sql = format!(
"SELECT table_name, fetched_at_epoch_ms FROM {} WHERE cache_key = {}",
naming::quote_ident(naming::CACHE_META_TABLE),
naming::quote_literal(key),
);
let df = self.cache.execute_sql(&sql)?;
if df.height() == 0 {
return Ok(None);
}
let table_name = as_str(df.column("table_name")?)?.value(0).to_string();
let fetched_at_epoch_ms = as_i64(df.column("fetched_at_epoch_ms")?)?.value(0);
Ok(Some(MemoEntry {
table_name,
fetched_at_epoch_ms,
}))
}
fn insert_memo(
&self,
key: &str,
sql: &str,
table: &str,
byte_estimate: i64,
row_count: i64,
) -> Result<()> {
let now = now_ms();
let stmt = format!(
"INSERT OR REPLACE INTO {} \
(cache_key, sql, table_name, fetched_at_epoch_ms, last_accessed_epoch_ms, \
byte_estimate, row_count) \
VALUES ({}, {}, {}, {}, {}, {}, {})",
naming::quote_ident(naming::CACHE_META_TABLE),
naming::quote_literal(key),
naming::quote_literal(sql),
naming::quote_literal(table),
now,
now,
byte_estimate,
row_count,
);
self.cache.execute_sql(&stmt)?;
Ok(())
}
fn touch(&self, key: &str) -> Result<()> {
let stmt = format!(
"UPDATE {} SET last_accessed_epoch_ms = {} WHERE cache_key = {}",
naming::quote_ident(naming::CACHE_META_TABLE),
now_ms(),
naming::quote_literal(key),
);
self.cache.execute_sql(&stmt)?;
Ok(())
}
fn drop_entry(&self, key: &str, table: &str) -> Result<()> {
self.cache.unregister(table)?;
let del = format!(
"DELETE FROM {} WHERE cache_key = {}",
naming::quote_ident(naming::CACHE_META_TABLE),
naming::quote_literal(key),
);
self.cache.execute_sql(&del)?;
Ok(())
}
fn evict_over_budget(&self) -> Result<()> {
let sum_sql = format!(
"SELECT CAST(COALESCE(SUM(byte_estimate), 0) AS BIGINT) AS n FROM {}",
naming::quote_ident(naming::CACHE_META_TABLE)
);
loop {
let df = self.cache.execute_sql(&sum_sql)?;
let total = if df.height() == 0 {
0
} else {
as_i64(df.column("n")?)?.value(0)
};
if total <= self.config.max_bytes as i64 {
return Ok(());
}
let pick = format!(
"SELECT cache_key, table_name FROM {} ORDER BY last_accessed_epoch_ms ASC LIMIT 1",
naming::quote_ident(naming::CACHE_META_TABLE)
);
let df = self.cache.execute_sql(&pick)?;
if df.height() == 0 {
return Ok(());
}
let key = as_str(df.column("cache_key")?)?.value(0).to_string();
let table = as_str(df.column("table_name")?)?.value(0).to_string();
self.drop_entry(&key, &table)?;
}
}
fn references_cache_resident(&self, sql: &str) -> bool {
if crate::parser::extract_builtin_dataset_names(sql)
.map(|d| !d.is_empty())
.unwrap_or(false)
{
return true;
}
let Ok(refs) = crate::parser::extract_table_refs(sql) else {
return false;
};
let resident = self.resident.borrow();
refs.iter()
.any(|t| t == naming::CACHE_META_TABLE || resident.contains(t))
}
fn clear_memo(&self) -> Result<()> {
self.ensure_meta_table()?;
let df = self.cache.execute_sql(&format!(
"SELECT cache_key, table_name FROM {}",
naming::quote_ident(naming::CACHE_META_TABLE)
))?;
let n = df.height();
if n == 0 {
return Ok(());
}
let keys = as_str(df.column("cache_key")?)?;
let tables = as_str(df.column("table_name")?)?;
let mut failures: Vec<String> = Vec::new();
for i in 0..n {
let key = keys.value(i);
let table = tables.value(i);
if let Err(e) = self.drop_entry(key, table) {
failures.push(format!("{key}: {e}"));
}
}
if !failures.is_empty() {
return Err(crate::GgsqlError::ReaderError(format!(
"clear_cache: {} cache entries failed to drop: {}",
failures.len(),
failures.join("; ")
)));
}
Ok(())
}
fn annotate_error(&self, context: &str, err: crate::GgsqlError) -> crate::GgsqlError {
let crate::GgsqlError::ReaderError(message) = &err else {
return err;
};
let inner = message
.strip_prefix("Failed to execute SQL: ")
.or_else(|| message.strip_prefix("Failed to prepare SQL: "))
.unwrap_or(message);
crate::GgsqlError::ReaderError(format!("{context}: {inner}"))
}
}
impl Reader for CachingReader {
fn execute_sql(&self, sql: &str) -> Result<DataFrame> {
self.ensure_meta_table()?;
if self.references_cache_resident(sql) {
return self.cache.execute_sql(sql);
}
if !self.config.enabled {
return self.primary.execute_sql(sql);
}
let key = self.cache_key(sql);
if let Some(entry) = self.lookup_memo(&key)? {
let age_ms = (now_ms() - entry.fetched_at_epoch_ms).max(0);
let ttl_ms = (self.config.ttl_secs as i64).saturating_mul(1000);
if age_ms < ttl_ms {
let select = format!("SELECT * FROM {}", naming::quote_ident(&entry.table_name));
if let Ok(df) = self.cache.execute_sql(&select) {
self.touch(&key)?;
return Ok(df);
}
}
self.drop_entry(&key, &entry.table_name)?;
}
let df = self.primary.execute_sql(sql)?;
if super::returns_rows(sql) && df.width() > 0 {
let table = naming::cache_result_table(&key);
let byte_estimate = estimate_bytes(&df);
let row_count = df.height() as i64;
self.cache.register(&table, df.clone(), true)?;
if let Err(e) = self.insert_memo(&key, sql, &table, byte_estimate, row_count) {
let _ = self.cache.unregister(&table);
return Err(e);
}
self.evict_over_budget()?;
}
Ok(df)
}
fn execute_sql_cached(&self, sql: &str) -> Result<DataFrame> {
self.cache.execute_sql(sql).map_err(|e| {
self.annotate_error(&format!("on the `{}` cache backend", self.cache_scheme), e)
})
}
fn register(&self, name: &str, df: DataFrame, replace: bool) -> Result<()> {
self.cache.register(name, df, replace)?;
self.resident.borrow_mut().insert(name.to_string());
Ok(())
}
fn unregister(&self, name: &str) -> Result<()> {
self.cache.unregister(name)?;
self.resident.borrow_mut().remove(name);
Ok(())
}
fn execute(&self, query: &str) -> Result<Spec> {
execute_with_reader(self, query)
}
fn dialect(&self) -> &dyn SqlDialect {
self.cache.dialect()
}
fn materialize_table(
&self,
name: &str,
column_aliases: &[String],
body_sql: &str,
) -> Result<()> {
let body = super::wrap_with_column_aliases(body_sql, column_aliases);
let df = self.execute_sql(&body)?;
self.register(name, df, true)
}
fn caches_sources(&self) -> bool {
true
}
fn clear_cache(&self) -> Result<()> {
self.clear_memo()
}
fn list_catalogs(&self) -> Result<Vec<String>> {
self.primary.list_catalogs()
}
fn list_schemas(&self, catalog: &str) -> Result<Vec<String>> {
self.primary.list_schemas(catalog)
}
fn list_tables(&self, catalog: &str, schema: &str) -> Result<Vec<TableInfo>> {
self.primary.list_tables(catalog, schema)
}
fn list_columns(&self, catalog: &str, schema: &str, table: &str) -> Result<Vec<ColumnInfo>> {
self.primary.list_columns(catalog, schema, table)
}
}
#[cfg(all(test, feature = "duckdb"))]
mod behavior_tests {
use super::*;
use crate::array_util::as_i64;
use crate::df;
use crate::reader::test_support::{ReadOnlyReader, SpyReader};
use crate::reader::{CacheBackend, DuckDBReader};
#[test]
fn test_register_writes_to_cache_and_query_routes_there() {
let (primary, log) = SpyReader::wrap(Box::new(DuckDBReader::new_in_memory().unwrap()));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
reader
.register("t", df! { "x" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let out = reader
.execute_sql_cached("SELECT COUNT(*) AS n FROM t")
.unwrap();
assert_eq!(as_i64(out.column("n").unwrap()).unwrap().value(0), 3);
assert!(log.lock().unwrap().is_empty());
}
#[test]
fn test_source_read_hits_primary_and_memoizes() {
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register("base", df! { "y" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let (primary, log) = SpyReader::wrap(Box::new(inner));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let q = "SELECT y FROM base ORDER BY y";
let d1 = reader.execute_sql(q).unwrap();
let d2 = reader.execute_sql(q).unwrap();
assert_eq!(d1.height(), 3);
assert_eq!(d2.height(), 3);
let hits = log
.lock()
.unwrap()
.iter()
.filter(|s| s.as_str() == q)
.count();
assert_eq!(hits, 1);
}
#[test]
fn test_full_execute_keeps_computation_off_primary() {
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register(
"sales",
df! { "x" => vec![1_i64, 2, 3, 4], "y" => vec![10_i64, 20, 30, 40] }.unwrap(),
true,
)
.unwrap();
let (primary, log) = SpyReader::wrap(Box::new(inner));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
reader
.execute("SELECT x, y FROM sales VISUALISE x, y DRAW point")
.unwrap();
let log = log.lock().unwrap();
assert!(
log.iter().all(|s| !s.to_uppercase().contains("TEMP TABLE")),
"primary must not be written to: {:?}",
*log
);
assert!(
log.iter().all(|s| !s.contains("__ggsql_")),
"primary must not see derived tables: {:?}",
*log
);
assert!(log.iter().any(|s| s.contains("sales")));
}
#[test]
fn test_caching_makes_read_only_primary_usable() {
let query = "SELECT v FROM t VISUALISE v AS x DRAW histogram";
let bare_primary = DuckDBReader::new_in_memory().unwrap();
bare_primary
.register(
"t",
df! { "v" => vec![1.0_f64, 2.0, 3.0, 4.0] }.unwrap(),
true,
)
.unwrap();
let bare = ReadOnlyReader::new(Box::new(bare_primary));
assert!(
bare.execute(query).is_err(),
"a read-only primary with no cache must fail to materialize"
);
let primary = DuckDBReader::new_in_memory().unwrap();
primary
.register(
"t",
df! { "v" => vec![1.0_f64, 2.0, 3.0, 4.0] }.unwrap(),
true,
)
.unwrap();
let cached = CachingReader::new(
Box::new(ReadOnlyReader::new(Box::new(primary))),
Box::new(DuckDBReader::new_in_memory().unwrap()),
"test://primary",
"duckdb",
);
assert!(
cached.execute(query).is_ok(),
"caching should make a read-only primary usable"
);
}
#[test]
fn test_no_cache_path_materializes_on_the_reader() {
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register(
"sales",
df! { "x" => vec![1_i64, 2, 3], "y" => vec![10_i64, 20, 30] }.unwrap(),
true,
)
.unwrap();
let (reader, log) = SpyReader::wrap(Box::new(inner));
reader
.execute("SELECT x, y FROM sales VISUALISE x, y DRAW point")
.unwrap();
assert!(
log.lock()
.unwrap()
.iter()
.any(|s| s.to_uppercase().contains("TEMP TABLE")),
"default path must materialize on the reader"
);
}
#[cfg(feature = "sqlite")]
#[test]
fn test_dialect_returns_cache_dialect() {
use crate::reader::SqliteReader;
let primary = Box::new(SqliteReader::new().unwrap());
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
assert_eq!(reader.dialect().sql_greatest(&["a", "b"]), "GREATEST(a, b)");
}
#[cfg(feature = "sqlite")]
#[test]
fn test_explicit_layer_source_with_stat_heterogeneous() {
use crate::reader::SqliteReader;
let primary = SqliteReader::new().unwrap();
primary
.register(
"tbl",
df! { "val" => vec![1.0_f64, 2.0, 2.0, 3.0, 3.0, 3.0, 9.0] }.unwrap(),
true,
)
.unwrap();
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(Box::new(primary), cache, "test://primary", "duckdb");
let spec = reader.execute("VISUALISE x DRAW histogram MAPPING val AS x FROM tbl");
assert!(
spec.is_ok(),
"explicit-source histogram failed: {:?}",
spec.err()
);
}
#[test]
fn test_aliased_cte_reading_primary_routes_to_primary() {
let base = DuckDBReader::new_in_memory().unwrap();
base.register("base", df! { "v" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let primary = Box::new(ReadOnlyReader::new(Box::new(base)));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(
"WITH t(a) AS (SELECT v FROM base) SELECT a FROM t VISUALISE a AS x DRAW point",
);
assert!(
spec.is_ok(),
"aliased CTE over a primary table should succeed: {:?}",
spec.err()
);
}
#[test]
fn test_aliased_cte_referencing_prior_cte_routes_to_cache() {
let base = DuckDBReader::new_in_memory().unwrap();
base.register("base", df! { "v" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let primary = Box::new(ReadOnlyReader::new(Box::new(base)));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(
"WITH a(p) AS (SELECT v FROM base), b(q) AS (SELECT p FROM a) \
SELECT q FROM b VISUALISE q AS x DRAW point",
);
assert!(
spec.is_ok(),
"dependent aliased CTE should succeed: {:?}",
spec.err()
);
}
#[test]
fn test_cte_joined_against_primary_base_table() {
let base = DuckDBReader::new_in_memory().unwrap();
base.register(
"base",
df! { "k" => vec![1_i64, 2], "w" => vec![10_i64, 20] }.unwrap(),
true,
)
.unwrap();
let primary = Box::new(ReadOnlyReader::new(Box::new(base)));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(
"WITH t AS (SELECT 1 AS k, 100 AS v) \
SELECT t.v, base.w FROM t JOIN base ON t.k = base.k \
VISUALISE v AS x, w AS y DRAW point",
);
assert!(
spec.is_ok(),
"CTE joined against a primary base table should succeed: {:?}",
spec.err()
);
}
#[test]
fn test_schema_qualified_base_table_join() {
let base = DuckDBReader::new_in_memory().unwrap();
base.execute_sql("CREATE SCHEMA myschema").unwrap();
base.execute_sql("CREATE TABLE myschema.base AS SELECT 1 AS k, 10 AS w")
.unwrap();
let primary = Box::new(ReadOnlyReader::new(Box::new(base)));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(
"WITH t AS (SELECT 1 AS k, 100 AS v) \
SELECT t.v, myschema.base.w FROM t JOIN myschema.base ON t.k = myschema.base.k \
VISUALISE v AS x, w AS y DRAW point",
);
assert!(
spec.is_ok(),
"CTE joined against a schema-qualified base table should succeed: {:?}",
spec.err()
);
}
#[test]
fn test_multiple_primary_joins_in_chain() {
let base = DuckDBReader::new_in_memory().unwrap();
base.execute_sql("CREATE TABLE base AS SELECT 1 AS k, 10 AS w")
.unwrap();
base.execute_sql("CREATE TABLE base2 AS SELECT 1 AS k, 20 AS z")
.unwrap();
let primary = Box::new(ReadOnlyReader::new(Box::new(base)));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(
"WITH t AS (SELECT 1 AS k) \
SELECT base.w, base2.z FROM t JOIN base ON t.k = base.k JOIN base2 ON t.k = base2.k \
VISUALISE w AS x, z AS y DRAW point",
);
assert!(
spec.is_ok(),
"two primaries in a join chain should both be staged: {:?}",
spec.err()
);
}
#[test]
fn test_same_named_schema_tables_joined() {
let base = DuckDBReader::new_in_memory().unwrap();
base.execute_sql("CREATE SCHEMA a").unwrap();
base.execute_sql("CREATE SCHEMA b").unwrap();
base.execute_sql("CREATE TABLE a.base AS SELECT 1 AS k, 10 AS w")
.unwrap();
base.execute_sql("CREATE TABLE b.base AS SELECT 1 AS k, 20 AS w")
.unwrap();
let primary = Box::new(ReadOnlyReader::new(Box::new(base)));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(
"WITH t AS (SELECT 1 AS k) \
SELECT a.base.w AS aw, b.base.w AS bw \
FROM t JOIN a.base ON t.k = a.base.k JOIN b.base ON t.k = b.base.k \
VISUALISE aw AS x, bw AS y DRAW point",
);
assert!(
spec.is_ok(),
"same-named schema tables should stage to distinct aliases: {:?}",
spec.err()
);
}
#[test]
fn test_setup_dml_runs_before_staging() {
let primary = Box::new(DuckDBReader::new_in_memory().unwrap());
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(
"CREATE TABLE sales AS SELECT * FROM (VALUES (1, 10), (2, 20)) t(k, w); \
WITH c AS (SELECT 1 AS k UNION ALL SELECT 2) \
SELECT sales.k, sales.w FROM sales JOIN c ON sales.k = c.k \
VISUALISE k AS x, w AS y DRAW point",
);
assert!(
spec.is_ok(),
"setup DML must run before staging reads the created table: {:?}",
spec.err()
);
}
#[test]
fn test_multi_statement_dml_ordering_before_staging() {
let primary = Box::new(DuckDBReader::new_in_memory().unwrap());
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let prepared = crate::execute::prepare_data_with_reader(
"CREATE TABLE t(k INTEGER, w INTEGER); \
INSERT INTO t VALUES (1, 10), (2, 20); \
UPDATE t SET w = w + 5 WHERE k = 1; \
WITH c AS (SELECT 1 AS k UNION ALL SELECT 2) \
SELECT t.k, t.w FROM t JOIN c ON t.k = c.k \
VISUALISE k AS x, w AS y DRAW point",
&reader,
)
.expect("multi-statement setup DML + staging should succeed");
let df = prepared.data.get(&naming::layer_key(0)).unwrap();
assert_eq!(df.height(), 2);
let ws = df.column("__ggsql_aes_pos2__").unwrap();
let mut vals: Vec<String> = (0..df.height())
.map(|i| crate::array_util::value_to_string(ws, i))
.collect();
vals.sort();
assert_eq!(vals, vec!["15".to_string(), "20".to_string()]);
}
#[test]
fn test_fully_quoted_schema_qualified_join() {
let base = DuckDBReader::new_in_memory().unwrap();
base.execute_sql("CREATE SCHEMA myschema").unwrap();
base.execute_sql("CREATE TABLE myschema.base AS SELECT 1 AS k, 10 AS w")
.unwrap();
let primary = Box::new(ReadOnlyReader::new(Box::new(base)));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(
"WITH t AS (SELECT 1 AS k) \
SELECT \"myschema\".\"base\".w FROM t JOIN \"myschema\".\"base\" \
ON t.k = \"myschema\".\"base\".k \
VISUALISE w AS x DRAW point",
);
assert!(
spec.is_ok(),
"fully-quoted schema-qualified join should succeed: {:?}",
spec.err()
);
}
#[test]
fn test_mixed_later_cte_body_joins_primary() {
let base = DuckDBReader::new_in_memory().unwrap();
base.register(
"base",
df! { "k" => vec![1_i64], "w" => vec![10_i64] }.unwrap(),
true,
)
.unwrap();
let primary = Box::new(ReadOnlyReader::new(Box::new(base)));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(
"WITH a AS (SELECT 1 AS k, 5 AS p), \
b AS (SELECT a.p, base.w FROM a JOIN base ON a.k = base.k) \
SELECT * FROM b VISUALISE p AS x, w AS y DRAW point",
);
assert!(
spec.is_ok(),
"mixed later-CTE body joining a primary table should succeed: {:?}",
spec.err()
);
}
#[test]
fn test_meta_table_queryable_before_any_read() {
let reader = CachingReader::new(
Box::new(DuckDBReader::new_in_memory().unwrap()),
Box::new(DuckDBReader::new_in_memory().unwrap()),
"test://primary",
"duckdb",
);
let meta = reader
.execute_sql(&format!(
"SELECT * FROM {}",
naming::quote_ident(naming::CACHE_META_TABLE)
))
.unwrap();
assert_eq!(meta.height(), 0);
}
#[test]
fn test_compute_surface_errors_name_the_cache_backend() {
let reader = CachingReader::new(
Box::new(DuckDBReader::new_in_memory().unwrap()),
Box::new(DuckDBReader::new_in_memory().unwrap()),
"test://primary",
"duckdb",
);
let err = reader
.execute_sql_cached("SELECT * FROM __ggsql_definitely_missing__")
.unwrap_err()
.to_string();
assert!(err.contains("`duckdb` cache backend"), "got: {err}");
assert!(err.contains("__ggsql_definitely_missing__"), "got: {err}");
assert!(!err.contains("Failed to prepare SQL"), "got: {err}");
assert!(!err.contains("Failed to execute SQL"), "got: {err}");
}
#[test]
fn test_meta_table_records_and_serves_memo() {
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register("base", df! { "y" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let (primary, log) = SpyReader::wrap(Box::new(inner));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let q = "SELECT y FROM base ORDER BY y";
reader.execute_sql(q).unwrap();
let meta = reader
.execute_sql(&format!("SELECT sql FROM {}", naming::CACHE_META_TABLE))
.unwrap();
assert_eq!(meta.height(), 1);
assert_eq!(
crate::array_util::as_str(meta.column("sql").unwrap())
.unwrap()
.value(0),
q
);
reader.execute_sql(q).unwrap();
let hits = log
.lock()
.unwrap()
.iter()
.filter(|s| s.as_str() == q)
.count();
assert_eq!(hits, 1);
}
#[cfg(feature = "sqlite")]
#[test]
fn test_sqlite_cache_backend_memoizes() {
use crate::reader::SqliteReader;
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register("base", df! { "y" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let (primary, log) = SpyReader::wrap(Box::new(inner));
let cache = Box::new(SqliteReader::new().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "sqlite");
let q = "SELECT y FROM base ORDER BY y";
let d1 = reader.execute_sql(q).unwrap();
let d2 = reader.execute_sql(q).unwrap();
assert_eq!(d1.height(), 3);
assert_eq!(d2.height(), 3);
let hits = log
.lock()
.unwrap()
.iter()
.filter(|s| s.as_str() == q)
.count();
assert_eq!(
hits, 1,
"SQLite cache should serve the repeat from the memo"
);
}
#[test]
fn test_resident_substring_not_false_matched() {
let primary = DuckDBReader::new_in_memory().unwrap();
primary
.register(
"orders_archive",
df! { "v" => vec![1_i64, 2, 3] }.unwrap(),
true,
)
.unwrap();
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(Box::new(primary), cache, "test://primary", "duckdb");
reader
.register("orders", df! { "v" => vec![9_i64] }.unwrap(), true)
.unwrap();
let df = reader.execute_sql("SELECT v FROM orders_archive").unwrap();
assert_eq!(df.height(), 3);
}
#[test]
fn test_default_config_enabled_ttl_300() {
let reader = CachingReader::with_config(
Box::new(DuckDBReader::new_in_memory().unwrap()),
Box::new(DuckDBReader::new_in_memory().unwrap()),
"test://primary",
"duckdb",
CacheConfig::default(),
);
assert!(reader.cache_config().enabled);
assert_eq!(reader.cache_config().ttl_secs, 300);
}
#[test]
fn test_repeat_query_hits_primary_once() {
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register("base", df! { "y" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let (primary, log) = SpyReader::wrap(Box::new(inner));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let q = "SELECT y FROM base ORDER BY y";
reader.execute_sql(q).unwrap();
reader.execute_sql(q).unwrap();
let hits = log
.lock()
.unwrap()
.iter()
.filter(|s| s.as_str() == q)
.count();
assert_eq!(hits, 1, "the repeat read is served from the memo");
}
#[test]
fn test_ttl_zero_always_misses() {
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register("base", df! { "y" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let (primary, log) = SpyReader::wrap(Box::new(inner));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::with_config(
primary,
cache,
"test://primary",
"duckdb",
CacheConfig {
enabled: true,
ttl_secs: 0,
max_bytes: 512 * 1024 * 1024,
},
);
let q = "SELECT y FROM base ORDER BY y";
reader.execute_sql(q).unwrap();
reader.execute_sql(q).unwrap();
let hits = log
.lock()
.unwrap()
.iter()
.filter(|s| s.as_str() == q)
.count();
assert_eq!(hits, 2, "ttl=0 must miss on every read");
}
#[test]
fn test_disabled_always_hits_primary() {
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register("base", df! { "y" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let (primary, log) = SpyReader::wrap(Box::new(inner));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::with_config(
primary,
cache,
"test://primary",
"duckdb",
CacheConfig {
enabled: false,
ttl_secs: 300,
max_bytes: 512 * 1024 * 1024,
},
);
let q = "SELECT y FROM base ORDER BY y";
reader.execute_sql(q).unwrap();
reader.execute_sql(q).unwrap();
let hits = log
.lock()
.unwrap()
.iter()
.filter(|s| s.as_str() == q)
.count();
assert_eq!(hits, 2, "a disabled cache always hits the primary");
let meta = reader
.execute_sql(&format!("SELECT * FROM {}", naming::CACHE_META_TABLE))
.unwrap();
assert_eq!(meta.height(), 0);
}
#[test]
fn test_lru_evicts_oldest_when_over_budget() {
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register(
"base",
df! { "a" => vec![1_i64, 2, 3], "b" => vec![4_i64, 5, 6] }.unwrap(),
true,
)
.unwrap();
let (primary, log) = SpyReader::wrap(Box::new(inner));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::with_config(
primary,
cache,
"test://primary",
"duckdb",
CacheConfig {
enabled: true,
ttl_secs: 300,
max_bytes: 1,
},
);
let q1 = "SELECT a FROM base ORDER BY a";
let q2 = "SELECT b FROM base ORDER BY b";
reader.execute_sql(q1).unwrap();
reader.execute_sql(q2).unwrap();
reader.execute_sql(q1).unwrap();
let q1_hits = log
.lock()
.unwrap()
.iter()
.filter(|s| s.as_str() == q1)
.count();
assert_eq!(
q1_hits, 2,
"the evicted entry is re-fetched from the primary"
);
let meta = reader
.execute_sql(&format!("SELECT * FROM {}", naming::CACHE_META_TABLE))
.unwrap();
assert!(meta.height() <= 1, "over-budget entries are evicted");
}
#[test]
fn test_missing_cached_table_self_heals() {
let inner = DuckDBReader::new_in_memory().unwrap();
inner
.register("base", df! { "y" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let (primary, log) = SpyReader::wrap(Box::new(inner));
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let q = "SELECT y FROM base ORDER BY y";
reader.execute_sql(q).unwrap();
let table = reader
.execute_sql(&format!(
"SELECT table_name FROM {}",
naming::CACHE_META_TABLE
))
.unwrap();
let table_name = crate::array_util::as_str(table.column("table_name").unwrap())
.unwrap()
.value(0)
.to_string();
reader
.cache
.execute_sql(&format!("DROP TABLE {}", naming::quote_ident(&table_name)))
.unwrap();
let df = reader.execute_sql(q).unwrap();
assert_eq!(df.height(), 3);
let hits = log
.lock()
.unwrap()
.iter()
.filter(|s| s.as_str() == q)
.count();
assert_eq!(hits, 2, "the missing table forced a primary re-fetch");
}
#[test]
fn test_parse_human_bytes() {
assert_eq!(parse_human_bytes("1024"), Some(1024));
assert_eq!(parse_human_bytes("512mb"), Some(512 * 1024 * 1024));
assert_eq!(parse_human_bytes("1GB"), Some(1024 * 1024 * 1024));
assert_eq!(parse_human_bytes(" 2kb "), Some(2 * 1024));
assert_eq!(parse_human_bytes("nonsense"), None);
}
#[test]
fn test_config_merge_uri_wins() {
let base = CacheConfig::default();
let merged = base.merge(CacheConfigOverride {
enabled: Some(false),
ttl_secs: Some(60),
max_bytes: None,
});
assert!(!merged.enabled);
assert_eq!(merged.ttl_secs, 60);
assert_eq!(merged.max_bytes, 512 * 1024 * 1024);
}
#[test]
fn test_pure_sql_reads_primary_not_cache() {
let primary = DuckDBReader::new_in_memory().unwrap();
primary
.register("t", df! { "v" => vec![1_i64, 2, 3] }.unwrap(), true)
.unwrap();
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(Box::new(primary), cache, "test://primary", "duckdb");
let df = reader.execute_sql("SELECT v FROM t").unwrap();
assert_eq!(df.height(), 3);
assert!(
reader.execute_sql_cached("SELECT v FROM t").is_err(),
"compute surface should not find the primary-only table"
);
}
#[test]
fn test_cache_resident_table_as_layer_source() {
let primary = Box::new(DuckDBReader::new_in_memory().unwrap());
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
reader
.register(
"only_in_cache",
df! { "val" => vec![1.0_f64, 2.0, 2.0, 3.0, 3.0, 3.0, 9.0] }.unwrap(),
true,
)
.unwrap();
let spec = reader.execute("VISUALISE x DRAW histogram MAPPING val AS x FROM only_in_cache");
assert!(
spec.is_ok(),
"cache-resident layer source should succeed: {:?}",
spec.err()
);
}
#[cfg(feature = "sqlite")]
#[test]
fn test_file_layer_source_staged_via_cache() {
use crate::reader::SqliteReader;
let dir = std::env::temp_dir();
let path = dir.join(format!("ggsql_cache_file_test_{}.csv", std::process::id()));
std::fs::write(&path, "val\n1.0\n2.0\n2.0\n3.0\n3.0\n3.0\n9.0\n").unwrap();
let path_str = path.to_str().unwrap().to_string();
let primary = Box::new(SqliteReader::new().unwrap());
let cache = Box::new(DuckDBReader::new_in_memory().unwrap());
let reader = CachingReader::new(primary, cache, "test://primary", "duckdb");
let spec = reader.execute(&format!(
"VISUALISE x DRAW histogram MAPPING val AS x FROM '{}'",
path_str
));
let _ = std::fs::remove_file(&path);
assert!(
spec.is_ok(),
"file layer source via cache should succeed: {:?}",
spec.err()
);
}
}