use std::sync::{Arc, PoisonError, RwLock, RwLockReadGuard, RwLockWriteGuard};
use rudb_bind::{Bound, Parameters};
use rudb_catalog::{Catalog, Entry, View};
use rudb_common::{Cancel, Error, Field, Memory, Result, Value};
use rudb_parse::ast::Ast;
use crate::config::Config;
use crate::connection::{Connection, single};
use crate::prepared::Prepared;
use crate::result::QueryResult;
use crate::settings::Settings;
const MEMORY: &str = ":memory:";
#[derive(Debug, Clone)]
pub struct Database {
shared: Shared,
}
#[derive(Debug, Clone)]
pub(crate) struct Shared {
inner: Arc<Inner>,
}
#[derive(Debug)]
struct Inner {
catalog: RwLock<Catalog>,
settings: Settings,
memory: Memory,
}
impl Default for Database {
fn default() -> Self {
Self::new()
}
}
impl Database {
#[must_use]
pub fn new() -> Self {
Self::with_config(Config::default())
}
#[must_use]
pub fn with_config(config: Config) -> Self {
let memory = Memory::new(config.memory_limit());
let settings = Settings::new(config);
let inner = Inner { catalog: RwLock::new(Catalog::new()), settings, memory };
Self { shared: Shared { inner: Arc::new(inner) } }
}
#[must_use]
pub fn config(&self) -> Config {
self.shared.inner.settings.config()
}
#[must_use]
pub fn opened_with(&self) -> Config {
self.shared.inner.settings.defaults()
}
pub fn setting(&self, name: &str) -> Result<String> {
self.shared.inner.settings.value(name)
}
#[must_use]
pub fn memory(&self) -> &Memory {
&self.shared.inner.memory
}
pub fn open(path: &str) -> Result<Self> {
Self::open_with(path, Config::default())
}
pub fn open_with(path: &str, config: Config) -> Result<Self> {
if path.is_empty() || path == MEMORY {
return Ok(Self::with_config(config));
}
Err(Error::not_implemented(format!(
"cannot open \"{path}\", because there is no storage format yet, see \
https://github.com/tamnd/rudb/issues/103"
)))
}
#[must_use]
pub fn connect(&self) -> Connection {
Connection::new(self.shared.clone())
}
pub fn prepare(&self, sql: &str) -> Result<Prepared> {
Prepared::new(self.shared.clone(), sql)
}
pub fn with_catalog<T>(&self, read: impl FnOnce(&Catalog) -> T) -> T {
read(&self.shared.read())
}
pub fn with_catalog_mut<T>(&self, write: impl FnOnce(&mut Catalog) -> T) -> T {
write(&mut self.shared.write())
}
pub fn create_table(&self, name: &str, columns: Vec<Field>) -> Result<()> {
let parts: Vec<&str> = name.split('.').collect();
let mut catalog = self.shared.write();
let resolved = catalog.resolve_for_create(&parts)?;
catalog.create_table(resolved, columns)
}
pub fn drop_table(&self, name: &str) -> Result<()> {
let parts: Vec<&str> = name.split('.').collect();
let mut catalog = self.shared.write();
let resolved = catalog.resolve(&parts)?;
catalog.drop_table(&resolved)
}
pub fn append(&self, name: &str, rows: &[Vec<Value>]) -> Result<()> {
let parts: Vec<&str> = name.split('.').collect();
let mut catalog = self.shared.write();
let resolved = catalog.resolve(&parts)?;
catalog.table_mut(&resolved)?.append_rows(rows)
}
pub fn table_len(&self, name: &str) -> Result<usize> {
let parts: Vec<&str> = name.split('.').collect();
let catalog = self.shared.read();
let resolved = catalog.resolve(&parts)?;
Ok(catalog.table(&resolved)?.rows().len())
}
#[must_use]
pub fn table_names(&self) -> Vec<String> {
self.shared.read().tables().map(|table| table.name().table.clone()).collect()
}
pub fn table_sql(&self, name: &str) -> Result<String> {
let parts: Vec<&str> = name.split('.').collect();
let catalog = self.shared.read();
let resolved = catalog.resolve(&parts)?;
let table = catalog.table(&resolved)?;
let columns: Vec<String> = table
.columns()
.iter()
.map(|field| {
let null = if field.not_null { " NOT NULL" } else { "" };
format!("{} {}{null}", field.name, field.ty)
})
.collect();
Ok(format!("CREATE TABLE {}({});", resolved.table, columns.join(", ")))
}
pub fn query(&self, sql: &str) -> Result<QueryResult> {
self.shared.query(sql, &self.shared.token())
}
pub fn execute(&self, sql: &str) -> Result<QueryResult> {
self.shared.execute(sql, &self.shared.token())
}
pub fn plan(&self, sql: &str) -> Result<String> {
self.shared.plan(sql)
}
pub fn value(&self, sql: &str) -> Result<Value> {
single(&self.query(sql)?)
}
}
impl Shared {
fn read(&self) -> RwLockReadGuard<'_, Catalog> {
self.inner.catalog.read().unwrap_or_else(PoisonError::into_inner)
}
fn write(&self) -> RwLockWriteGuard<'_, Catalog> {
self.inner.catalog.write().unwrap_or_else(PoisonError::into_inner)
}
pub(crate) fn query(&self, sql: &str, cancel: &Cancel) -> Result<QueryResult> {
let catalog = self.read();
let plan = planned(sql, &catalog, &self.optimizer()?)?;
run(&plan, &catalog, cancel, &self.inner.memory)
}
fn optimizer(&self) -> Result<rudb_opt::pass::Context> {
rudb_opt::pass::Context::without(&self.inner.settings.disabled_optimizers())
}
pub(crate) fn timeout(&self) -> Option<std::time::Duration> {
self.inner.settings.config().query_timeout()
}
pub(crate) fn token(&self) -> Cancel {
match self.inner.settings.config().query_timeout() {
Some(timeout) => Cancel::after(timeout),
None => Cancel::new(),
}
}
pub(crate) fn plan(&self, sql: &str) -> Result<String> {
Ok(planned(sql, &self.read(), &self.optimizer()?)?.to_string())
}
pub(crate) fn execute(&self, sql: &str, cancel: &Cancel) -> Result<QueryResult> {
let ast = rudb_parse::parse_ast(sql)?;
self.execute_ast(&ast, &Parameters::new(), cancel)
}
pub(crate) fn execute_ast(
&self,
ast: &Ast,
parameters: &Parameters,
cancel: &Cancel,
) -> Result<QueryResult> {
let mut catalog = self.write();
let context = self.optimizer()?;
match rudb_bind::bind_statement_with(ast, &catalog, parameters)? {
Bound::Query(mut plan) => {
rudb_opt::optimize_with(&mut plan, &context)?;
run(&plan, &catalog, cancel, &self.inner.memory)
}
Bound::Setting(setting) => {
let value = setting.value.as_ref();
self.inner.settings.apply(
&self.inner.memory,
&setting.name,
setting.scope,
value,
)?;
Ok(QueryResult::empty())
}
Bound::CreateTable(create) => {
create_table(create, &mut catalog, cancel, &self.inner.memory, &context)?;
Ok(QueryResult::empty())
}
Bound::CreateView(create) => {
create_view(create, &mut catalog)?;
Ok(QueryResult::empty())
}
Bound::DropTable(drop) => {
for name in &drop.names {
match drop.kind {
Entry::Table => catalog.drop_table(name)?,
Entry::View => catalog.drop_view(name)?,
}
}
Ok(QueryResult::empty())
}
Bound::Insert(mut insert) => {
rudb_opt::optimize_with(&mut insert.source, &context)?;
let result = run(&insert.source, &catalog, cancel, &self.inner.memory)?;
let table = catalog.table_mut(&insert.name)?;
for chunk in result.into_chunks() {
table.append(chunk)?;
}
Ok(QueryResult::empty())
}
}
}
}
fn planned(
sql: &str,
catalog: &Catalog,
context: &rudb_opt::pass::Context,
) -> Result<rudb_plan::Plan> {
let mut plan = rudb_bind::bind_sql(sql, catalog)?;
rudb_opt::optimize_with(&mut plan, context)?;
Ok(plan)
}
fn run(
plan: &rudb_plan::Plan,
catalog: &Catalog,
cancel: &Cancel,
memory: &Memory,
) -> Result<QueryResult> {
let mut root = rudb_exec::build_with(plan, catalog, cancel, memory)?;
let names = root.schema().names();
let types = root.schema().types();
let mut held = memory.reservation();
let mut chunks = Vec::new();
while let Some(chunk) = root.next()? {
if chunk.is_empty() {
continue;
}
let chunk = chunk.flatten()?;
held.grow(u64::try_from(chunk.footprint()).unwrap_or(u64::MAX))?;
chunks.push(chunk);
}
Ok(QueryResult::new(names, types, chunks, held))
}
fn create_view(create: rudb_bind::CreateView, catalog: &mut Catalog) -> Result<()> {
if create.if_not_exists && catalog.entry(&create.name).is_ok() {
return Ok(());
}
if create.or_replace && catalog.view(&create.name).is_ok() {
catalog.drop_view(&create.name)?;
}
catalog.create_view(View::new(create.name, create.sql, create.aliases))
}
fn create_table(
mut create: rudb_bind::CreateTable,
catalog: &mut Catalog,
cancel: &Cancel,
memory: &Memory,
context: &rudb_opt::pass::Context,
) -> Result<()> {
if create.if_not_exists && catalog.table(&create.name).is_ok() {
return Ok(());
}
let rows = match &mut create.source {
Some(plan) => {
rudb_opt::optimize_with(plan, context)?;
Some(run(plan, catalog, cancel, memory)?)
}
None => None,
};
if create.or_replace && catalog.table(&create.name).is_ok() {
catalog.drop_table(&create.name)?;
}
catalog.create_table(create.name.clone(), create.columns)?;
if let Some(rows) = rows {
let table = catalog.table_mut(&create.name)?;
for chunk in rows.into_chunks() {
table.append(chunk)?;
}
}
Ok(())
}