use std::sync::{Arc, PoisonError, RwLock, RwLockReadGuard, RwLockWriteGuard};
use rudb_bind::{Bound, Parameters};
use rudb_catalog::Catalog;
use rudb_common::{Error, Field, Result, Value};
use rudb_parse::ast::Ast;
use crate::connection::{Connection, single};
use crate::prepared::Prepared;
use crate::result::QueryResult;
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>,
}
impl Default for Database {
fn default() -> Self {
Self::new()
}
}
impl Database {
#[must_use]
pub fn new() -> Self {
let inner = Inner { catalog: RwLock::new(Catalog::new()) };
Self { shared: Shared { inner: Arc::new(inner) } }
}
pub fn open(path: &str) -> Result<Self> {
if path.is_empty() || path == MEMORY {
return Ok(Self::new());
}
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)
}
pub fn execute(&self, sql: &str) -> Result<QueryResult> {
self.shared.execute(sql)
}
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) -> Result<QueryResult> {
let catalog = self.read();
let plan = planned(sql, &catalog)?;
run(&plan, &catalog)
}
pub(crate) fn plan(&self, sql: &str) -> Result<String> {
Ok(planned(sql, &self.read())?.to_string())
}
pub(crate) fn execute(&self, sql: &str) -> Result<QueryResult> {
let ast = rudb_parse::parse_ast(sql)?;
self.execute_ast(&ast, &Parameters::new())
}
pub(crate) fn execute_ast(&self, ast: &Ast, parameters: &Parameters) -> Result<QueryResult> {
let mut catalog = self.write();
match rudb_bind::bind_statement_with(ast, &catalog, parameters)? {
Bound::Query(mut plan) => {
rudb_opt::optimize(&mut plan)?;
run(&plan, &catalog)
}
Bound::CreateTable(create) => {
create_table(create, &mut catalog)?;
Ok(QueryResult::empty())
}
Bound::DropTable(drop) => {
for name in &drop.names {
catalog.drop_table(name)?;
}
Ok(QueryResult::empty())
}
Bound::Insert(mut insert) => {
rudb_opt::optimize(&mut insert.source)?;
let result = run(&insert.source, &catalog)?;
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) -> Result<rudb_plan::Plan> {
let mut plan = rudb_bind::bind_sql(sql, catalog)?;
rudb_opt::optimize(&mut plan)?;
Ok(plan)
}
fn run(plan: &rudb_plan::Plan, catalog: &Catalog) -> Result<QueryResult> {
let mut root = rudb_exec::build(plan, catalog)?;
let names = root.schema().names();
let types = root.schema().types();
let mut chunks = Vec::new();
while let Some(chunk) = root.next()? {
if chunk.is_empty() {
continue;
}
chunks.push(chunk.flatten()?);
}
Ok(QueryResult::new(names, types, chunks))
}
fn create_table(mut create: rudb_bind::CreateTable, catalog: &mut Catalog) -> 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(plan)?;
Some(run(plan, catalog)?)
}
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(())
}