use std::{
ffi::{c_void, CStr, CString},
mem,
os::raw::c_char,
ptr, str,
sync::Arc,
};
use crate::{
config::Config,
error::{Error, Result},
ffi::{
duckdb_close, duckdb_connect, duckdb_connection, duckdb_database, duckdb_disconnect,
duckdb_free, duckdb_open_ext, duckdb_query, duckdb_result, DuckDBError, DuckDBSuccess,
Error as FFIError,
},
helpers::duck_result::result_from_duckdb_result,
raw::{
appender::Appender,
result::DuckResult,
statement::{CachedStatement, Statement},
},
types::appendable::AppendAble,
};
pub struct RawDatabase(pub(crate) duckdb_database);
impl RawDatabase {
#[inline]
pub unsafe fn new(db: duckdb_database) -> Result<RawDatabase> {
if db.is_null() {
return Err(Error::DuckDBFailure(
FFIError::new(DuckDBError),
Some("database is null".to_owned()),
));
}
Ok(RawDatabase(db))
}
pub(crate) fn open_with_flags(
c_path: &CStr,
config: Config,
) -> Result<RawDatabase> {
unsafe {
let mut db: duckdb_database = ptr::null_mut();
let mut c_err = std::ptr::null_mut();
let r = duckdb_open_ext(c_path.as_ptr(), &mut db, config.duckdb_config(), &mut c_err);
if r != DuckDBSuccess {
let msg = Some(CStr::from_ptr(c_err).to_string_lossy().to_string());
duckdb_free(c_err as *mut c_void);
return Err(Error::DuckDBFailure(FFIError::new(r), msg));
}
RawDatabase::new(db)
}
}
}
unsafe impl Send for RawDatabase {}
unsafe impl Sync for RawDatabase {}
impl Drop for RawDatabase {
#[inline]
fn drop(&mut self) {
unsafe {
if !self.0.is_null() {
duckdb_close(&mut self.0);
}
}
}
}
pub struct RawConnection {
pub db: Arc<RawDatabase>,
pub con: duckdb_connection,
}
impl RawConnection {
#[allow(unused)]
fn raw(&self) -> duckdb_connection {
self.con
}
#[inline]
pub(crate) fn new(db: Arc<RawDatabase>) -> Result<RawConnection> {
let mut con: duckdb_connection = ptr::null_mut();
let r = unsafe { duckdb_connect(db.0, &mut con) };
if r != DuckDBSuccess {
unsafe { duckdb_disconnect(&mut con) };
return Err(Error::DuckDBFailure(FFIError::new(r), Some("connect error".to_owned())));
}
Ok(RawConnection { db, con })
}
pub fn open_with_flags(
c_path: &CStr,
config: Config,
) -> Result<RawConnection> {
RawConnection::new(Arc::new(RawDatabase::open_with_flags(c_path, config)?))
}
pub fn close(&mut self) -> Result<()> {
if self.con.is_null() {
return Ok(());
}
unsafe {
duckdb_disconnect(&mut self.con);
self.con = ptr::null_mut();
}
Ok(())
}
pub fn try_clone(&self) -> Result<Self> {
RawConnection::new(self.db.clone())
}
#[must_use = "query returns a DuckResult; discard explicitly with `let _ = ...` if not needed"]
pub fn query(
&mut self,
sql: impl AsRef<str>,
) -> Result<DuckResult> {
let c_str = CString::new(sql.as_ref())?;
let mut out = unsafe { mem::zeroed::<duckdb_result>() };
let r = unsafe {
duckdb_query(self.con, c_str.as_ptr() as *const c_char, &mut out as *mut duckdb_result)
};
result_from_duckdb_result(r, &mut out as *mut duckdb_result)?;
Ok(DuckResult::new(out))
}
#[must_use = "prepare returns a Statement; call execute() to run it"]
#[allow(unused)]
pub fn prepare(
&self,
sql: impl AsRef<str>,
) -> Result<Statement<'_>> {
Statement::new(self, sql.as_ref())
}
#[must_use = "appender returns an Appender that must be used to insert rows"]
pub fn appender(
&mut self,
table: &str,
schema: &str,
) -> Result<Appender> {
Appender::new(self.clone(), table, schema)
}
#[must_use = "insert result should be checked"]
#[allow(unused)]
pub fn insert<T: AppendAble, I>(
&mut self,
sql: &str,
values: I,
) -> Result<()>
where
I: IntoIterator<Item = T>,
{
let mut stmt = Statement::new(self, sql)?;
for mut each in values {
stmt.bind(&mut each)?;
}
let mut res = stmt.execute()?;
if res.changes() > 0 {
Ok(())
} else {
Err(Error::DuckDBFailure(
FFIError::new(DuckDBError),
Some("Failed to insert values".to_owned()),
))
}
}
#[must_use = "the DuckResult carries both affected-row count (.changes()) and row iterator"]
pub fn execute(
&mut self,
sql: impl AsRef<str>,
binds: &mut [&mut dyn AppendAble],
) -> Result<DuckResult> {
let mut stmt = CachedStatement::prepare(self, sql)?;
for (i, bind) in binds.iter_mut().enumerate() {
stmt.bind((i + 1) as u64, *bind)?;
}
stmt.execute()
}
}
impl Clone for RawConnection {
fn clone(&self) -> Self {
match self.try_clone() {
Ok(con) => con,
Err(e) => panic!("Failed to clone RawConnection: {e:?}"),
}
}
}
impl Drop for RawConnection {
#[inline]
fn drop(&mut self) {
use std::thread::panicking;
if let Err(e) = self.close() {
if panicking() {
eprintln!("Error while closing DuckDB connection: {e:?}");
} else {
panic!("Error while closing DuckDB connection: {e:?}");
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_raw_connection_open() {
let path = CString::new(":memory:").unwrap();
let config = Config::default();
let conn = RawConnection::open_with_flags(&path, config);
assert!(conn.is_ok());
}
#[test]
fn test_raw_connection_execute() {
let path = CString::new(":memory:").unwrap();
let config = Config::default();
let mut conn = RawConnection::open_with_flags(&path, config).unwrap();
let result = conn.query("CREATE TABLE test (id INTEGER PRIMARY KEY, name TEXT)");
assert!(result.is_ok(), "{}", result.err().unwrap());
}
#[test]
fn test_raw_connection_prepare() {
let path = CString::new(":memory:").unwrap();
let config = Config::default();
let mut conn = RawConnection::open_with_flags(&path, config).unwrap();
let result = conn.query("CREATE TABLE test (id INTEGER PRIMARY KEY, name TEXT)");
assert!(result.is_ok(), "{}", result.err().unwrap());
let stmt = conn.prepare("SELECT * FROM test");
assert!(stmt.is_ok(), "{}", stmt.err().unwrap());
}
#[test]
fn test_raw_connection_appender() {
let path = CString::new(":memory:").unwrap();
let config = Config::default();
let mut conn = RawConnection::open_with_flags(&path, config).unwrap();
let result = conn.query("CREATE TABLE test_table (id INTEGER PRIMARY KEY, name TEXT)");
assert!(result.is_ok(), "{}", result.err().unwrap());
let appender = conn.appender("test_table", "main");
assert!(appender.is_ok(), "{}", appender.err().unwrap());
}
}