use crate::error::{QueryError, Result};
use crate::types::QueryResult;
use manifold::column_family::ColumnFamilyDatabase;
use manifold_graph::{GraphTable, GraphTableRead};
use manifold_properties::{PropertyTable, PropertyTableRead, PropertyValue};
use manifold_vectors::{VectorTable, VectorTableRead};
use std::path::{Path, PathBuf};
use uuid::Uuid;
pub struct Database {
path: PathBuf,
cf_db: ColumnFamilyDatabase,
}
impl Database {
pub async fn open<P: AsRef<Path>>(path: P) -> Result<Self> {
let path = path.as_ref().to_path_buf();
if let Some(parent) = path.parent() {
if !parent.exists() {
std::fs::create_dir_all(parent).map_err(|e| QueryError::ConnectionError {
message: format!("Failed to create parent directory: {}", e),
})?;
}
}
let cf_db = ColumnFamilyDatabase::builder().open(&path).map_err(|e| {
QueryError::ConnectionError {
message: format!("Failed to open Manifold database: {}", e),
}
})?;
Ok(Self { path, cf_db })
}
pub async fn close(self) -> Result<()> {
drop(self.cf_db);
Ok(())
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn collection(&self, name: &str) -> Result<Collection> {
let cf =
self.cf_db
.column_family_or_create(name)
.map_err(|e| QueryError::ConnectionError {
message: format!("Failed to get column family '{}': {}", name, e),
})?;
Ok(Collection { cf })
}
pub async fn execute_hyperql(&self, _query: &str) -> Result<QueryResult> {
Ok(QueryResult {
rows: Vec::new(),
affected_rows: 0,
})
}
pub async fn execute_sql(&self, _query: &str) -> Result<QueryResult> {
Ok(QueryResult {
rows: Vec::new(),
affected_rows: 0,
})
}
pub async fn execute_cypher(&self, _query: &str) -> Result<QueryResult> {
Ok(QueryResult {
rows: Vec::new(),
affected_rows: 0,
})
}
pub async fn execute_custom(&self, language: &str, query: &str) -> Result<QueryResult> {
match language {
"hyperql" => self.execute_hyperql(query).await,
"sql" => self.execute_sql(query).await,
"cypher" => self.execute_cypher(query).await,
_ => Err(QueryError::UnsupportedLanguage {
language: language.to_string(),
}),
}
}
}
pub struct Collection {
cf: manifold::column_family::ColumnFamily,
}
impl Collection {
pub fn create_entity(&self, id: Uuid, data: serde_json::Value) -> Result<()> {
let write_txn = self
.cf
.begin_write()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin transaction: {}", e),
})?;
let mut props = PropertyTable::open(&write_txn, "properties").map_err(|e| {
QueryError::ExecutionError {
message: format!("Failed to open PropertyTable: {}", e),
}
})?;
if let serde_json::Value::Object(map) = data {
for (key, value) in map {
let prop_value = match value {
serde_json::Value::Number(n) if n.is_i64() => {
PropertyValue::new_integer(n.as_i64().unwrap())
}
serde_json::Value::Number(n) if n.is_f64() => {
PropertyValue::new_float(n.as_f64().unwrap())
}
serde_json::Value::Bool(b) => PropertyValue::new_boolean(b),
serde_json::Value::String(s) => PropertyValue::new_string(s),
serde_json::Value::Null => PropertyValue::new_null(),
_ => PropertyValue::new_string(value.to_string()),
};
props.set(&id, key.as_str(), prop_value).map_err(|e| {
QueryError::ExecutionError {
message: format!("Failed to set property: {}", e),
}
})?;
}
}
drop(props);
write_txn.commit().map_err(|e| QueryError::ExecutionError {
message: format!("Failed to commit transaction: {}", e),
})?;
Ok(())
}
pub fn get_entity(&self, id: Uuid) -> Result<Option<serde_json::Value>> {
let read_txn = self
.cf
.begin_read()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin read transaction: {}", e),
})?;
let props = PropertyTableRead::open(&read_txn, "properties").map_err(|e| {
QueryError::ExecutionError {
message: format!("Failed to open PropertyTable: {}", e),
}
})?;
let mut map = serde_json::Map::new();
let properties = props.get_all(&id).map_err(|e| QueryError::ExecutionError {
message: format!("Failed to read properties: {}", e),
})?;
for (key, value_guard) in properties {
let value = value_guard.value();
let json_value = if let Some(i) = value.as_integer() {
serde_json::Value::Number(i.into())
} else if let Some(f) = value.as_float() {
serde_json::Number::from_f64(f)
.map(serde_json::Value::Number)
.unwrap_or(serde_json::Value::Null)
} else if let Some(b) = value.as_boolean() {
serde_json::Value::Bool(b)
} else if let Some(s) = value.as_string() {
serde_json::from_str(s).unwrap_or_else(|_| serde_json::Value::String(s.to_string()))
} else {
serde_json::Value::Null
};
map.insert(key.to_string(), json_value);
}
if map.is_empty() {
Ok(None)
} else {
Ok(Some(serde_json::Value::Object(map)))
}
}
pub fn update_entity(&self, id: Uuid, data: serde_json::Value) -> Result<()> {
self.create_entity(id, data)
}
pub fn delete_entity(&self, id: Uuid) -> Result<()> {
let write_txn = self
.cf
.begin_write()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin transaction: {}", e),
})?;
let mut props = PropertyTable::open(&write_txn, "properties").map_err(|e| {
QueryError::ExecutionError {
message: format!("Failed to open PropertyTable: {}", e),
}
})?;
let read_txn = self
.cf
.begin_read()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin read transaction: {}", e),
})?;
let read_props = PropertyTableRead::open(&read_txn, "properties").map_err(|e| {
QueryError::ExecutionError {
message: format!("Failed to open PropertyTable: {}", e),
}
})?;
let keys_to_delete: Vec<String> = read_props
.get_all(&id)
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to read properties: {}", e),
})?
.iter()
.map(|(key, _)| key.to_string())
.collect();
drop(read_props);
drop(read_txn);
let keys_refs: Vec<(Uuid, &str)> = keys_to_delete
.iter()
.map(|key| (id, key.as_str()))
.collect();
props
.remove_bulk(&keys_refs)
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to delete properties: {}", e),
})?;
drop(props);
write_txn.commit().map_err(|e| QueryError::ExecutionError {
message: format!("Failed to commit transaction: {}", e),
})?;
Ok(())
}
pub fn add_edge(&self, source: Uuid, edge_type: &str, target: Uuid) -> Result<()> {
let write_txn = self
.cf
.begin_write()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin transaction: {}", e),
})?;
let mut graph =
GraphTable::open(&write_txn, "edges").map_err(|e| QueryError::ExecutionError {
message: format!("Failed to open GraphTable: {}", e),
})?;
graph
.add_edge(&source, edge_type, &target, true, 1.0, None)
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to add edge: {}", e),
})?;
drop(graph);
write_txn.commit().map_err(|e| QueryError::ExecutionError {
message: format!("Failed to commit transaction: {}", e),
})?;
Ok(())
}
pub fn get_outgoing_edges(&self, source: Uuid, edge_type: &str) -> Result<Vec<Uuid>> {
let read_txn = self
.cf
.begin_read()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin read transaction: {}", e),
})?;
let graph =
GraphTableRead::open(&read_txn, "edges").map_err(|e| QueryError::ExecutionError {
message: format!("Failed to open GraphTable: {}", e),
})?;
let mut targets = Vec::new();
let edges = graph
.outgoing_edges(&source)
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to read outgoing edges: {}", e),
})?;
for edge_result in edges {
let edge = edge_result.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to read edge: {}", e),
})?;
if edge.edge_type == edge_type && edge.is_active {
targets.push(edge.target);
}
}
Ok(targets)
}
pub fn get_incoming_edges(&self, target: Uuid, edge_type: &str) -> Result<Vec<Uuid>> {
let read_txn = self
.cf
.begin_read()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin read transaction: {}", e),
})?;
let graph =
GraphTableRead::open(&read_txn, "edges").map_err(|e| QueryError::ExecutionError {
message: format!("Failed to open GraphTable: {}", e),
})?;
let mut sources = Vec::new();
let edges = graph
.incoming_edges(&target)
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to read incoming edges: {}", e),
})?;
for edge_result in edges {
let edge = edge_result.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to read edge: {}", e),
})?;
if edge.edge_type == edge_type && edge.is_active {
sources.push(edge.source);
}
}
Ok(sources)
}
pub fn list_all_ids(&self) -> Result<Vec<Uuid>> {
let read_txn = self
.cf
.begin_read()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin read transaction: {}", e),
})?;
let props = PropertyTableRead::open(&read_txn, "properties").map_err(|e| {
QueryError::ExecutionError {
message: format!("Failed to open PropertyTable: {}", e),
}
})?;
let mut entity_ids = std::collections::HashSet::new();
for result in props.iter().map_err(|e| QueryError::ExecutionError {
message: format!("Failed to iterate properties: {}", e),
})? {
let ((entity_id, _key), _value) = result.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to read property: {}", e),
})?;
entity_ids.insert(entity_id);
}
Ok(entity_ids.into_iter().collect())
}
pub fn list_all_entities(&self) -> Result<Vec<(Uuid, serde_json::Value)>> {
let ids = self.list_all_ids()?;
let mut entities = Vec::new();
for id in ids {
if let Some(data) = self.get_entity(id)? {
entities.push((id, data));
}
}
Ok(entities)
}
pub fn vectors<const DIM: usize>(&self, table_name: &str) -> Result<VectorTableWrapper<DIM>> {
Ok(VectorTableWrapper {
cf: self.cf.clone(),
table_name: table_name.to_string(),
})
}
}
pub struct VectorTableWrapper<const DIM: usize> {
cf: manifold::column_family::ColumnFamily,
table_name: String,
}
impl<const DIM: usize> VectorTableWrapper<DIM> {
pub fn insert(&self, id: &Uuid, vector: &[f32]) -> Result<()> {
if vector.len() != DIM {
return Err(QueryError::ExecutionError {
message: format!(
"Vector dimension mismatch: expected {}, got {}",
DIM,
vector.len()
),
});
}
let write_txn = self
.cf
.begin_write()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin transaction: {}", e),
})?;
let mut vectors = VectorTable::<DIM>::open(&write_txn, &self.table_name).map_err(|e| {
QueryError::ExecutionError {
message: format!("Failed to open VectorTable: {}", e),
}
})?;
let mut arr = [0.0f32; DIM];
arr.copy_from_slice(vector);
vectors
.insert(id, &arr)
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to insert vector: {}", e),
})?;
drop(vectors);
write_txn.commit().map_err(|e| QueryError::ExecutionError {
message: format!("Failed to commit transaction: {}", e),
})?;
Ok(())
}
pub fn get(&self, id: &Uuid) -> Result<Option<Vec<f32>>> {
let read_txn = self
.cf
.begin_read()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin read transaction: {}", e),
})?;
let vectors = VectorTableRead::<DIM>::open(&read_txn, &self.table_name).map_err(|e| {
QueryError::ExecutionError {
message: format!("Failed to open VectorTable: {}", e),
}
})?;
let result = vectors.get(id).map_err(|e| QueryError::ExecutionError {
message: format!("Failed to read vector: {}", e),
})?;
Ok(result.map(|guard| guard.value().to_vec()))
}
pub fn search_similar(&self, query: &[f32], limit: usize) -> Result<Vec<Uuid>> {
if query.len() != DIM {
return Err(QueryError::ExecutionError {
message: format!(
"Query vector dimension mismatch: expected {}, got {}",
DIM,
query.len()
),
});
}
let read_txn = self
.cf
.begin_read()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to begin read transaction: {}", e),
})?;
let vectors = VectorTableRead::<DIM>::open(&read_txn, &self.table_name).map_err(|e| {
QueryError::ExecutionError {
message: format!("Failed to open VectorTable: {}", e),
}
})?;
let mut similarities: Vec<(Uuid, f32)> = Vec::new();
let iter = vectors
.all_vectors()
.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to iterate vectors: {}", e),
})?;
for result in iter {
let (id, vector_guard) = result.map_err(|e| QueryError::ExecutionError {
message: format!("Failed to read vector entry: {}", e),
})?;
let vector = vector_guard.value();
let similarity = cosine_similarity(query, vector);
similarities.push((id, similarity));
}
similarities.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
let results: Vec<Uuid> = similarities
.into_iter()
.take(limit)
.map(|(id, _)| id)
.collect();
Ok(results)
}
}
fn cosine_similarity(a: &[f32], b: &[f32]) -> f32 {
debug_assert_eq!(a.len(), b.len(), "Vectors must have same length");
let mut dot_product = 0.0;
let mut norm_a = 0.0;
let mut norm_b = 0.0;
for i in 0..a.len() {
dot_product += a[i] * b[i];
norm_a += a[i] * a[i];
norm_b += b[i] * b[i];
}
let magnitude = (norm_a * norm_b).sqrt();
if magnitude == 0.0 {
0.0
} else {
dot_product / magnitude
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[tokio::test]
async fn test_database_open() {
let temp_dir = std::env::temp_dir().join("audb_test_manifold_open.manifold");
let wal_path = temp_dir.with_extension("wal");
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
let db = Database::open(&temp_dir).await;
assert!(db.is_ok());
let db = db.unwrap();
assert_eq!(db.path(), temp_dir.as_path());
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
}
#[tokio::test]
async fn test_create_and_get_entity() {
let temp_dir = std::env::temp_dir().join("audb_test_manifold_crud.manifold");
let wal_path = temp_dir.with_extension("wal");
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
let db = Database::open(&temp_dir).await.unwrap();
let collection = db.collection("users").unwrap();
let id = Uuid::new_v4();
let data = json!({
"name": "Alice",
"age": 42,
"active": true
});
collection.create_entity(id, data.clone()).unwrap();
let retrieved = collection.get_entity(id).unwrap();
assert!(retrieved.is_some());
let retrieved = retrieved.unwrap();
assert_eq!(retrieved["name"], "Alice");
assert_eq!(retrieved["age"], 42);
assert_eq!(retrieved["active"], true);
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
}
#[tokio::test]
async fn test_update_entity() {
let temp_dir = std::env::temp_dir().join("audb_test_manifold_update.manifold");
let wal_path = temp_dir.with_extension("wal");
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
let db = Database::open(&temp_dir).await.unwrap();
let collection = db.collection("users").unwrap();
let id = Uuid::new_v4();
let data = json!({ "name": "Alice", "age": 42 });
collection.create_entity(id, data).unwrap();
let updated = json!({ "name": "Alice", "age": 43 });
collection.update_entity(id, updated).unwrap();
let retrieved = collection.get_entity(id).unwrap().unwrap();
assert_eq!(retrieved["age"], 43);
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
}
#[tokio::test]
async fn test_delete_entity() {
let temp_dir = std::env::temp_dir().join("audb_test_manifold_delete.manifold");
let wal_path = temp_dir.with_extension("wal");
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
let db = Database::open(&temp_dir).await.unwrap();
let collection = db.collection("users").unwrap();
let id = Uuid::new_v4();
let data = json!({ "name": "Bob" });
collection.create_entity(id, data).unwrap();
assert!(collection.get_entity(id).unwrap().is_some());
collection.delete_entity(id).unwrap();
assert!(collection.get_entity(id).unwrap().is_none());
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
}
#[tokio::test]
async fn test_edges() {
let temp_dir = std::env::temp_dir().join("audb_test_manifold_edges.manifold");
let wal_path = temp_dir.with_extension("wal");
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
let db = Database::open(&temp_dir).await.unwrap();
let collection = db.collection("test").unwrap();
let user = Uuid::new_v4();
let post1 = Uuid::new_v4();
let post2 = Uuid::new_v4();
collection.add_edge(post1, "authored_by", user).unwrap();
collection.add_edge(post2, "authored_by", user).unwrap();
let authors1 = collection.get_outgoing_edges(post1, "authored_by").unwrap();
assert_eq!(authors1, vec![user]);
let authored_posts = collection.get_incoming_edges(user, "authored_by").unwrap();
assert_eq!(authored_posts.len(), 2);
assert!(authored_posts.contains(&post1));
assert!(authored_posts.contains(&post2));
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
}
#[tokio::test]
async fn test_vectors() {
let temp_dir = std::env::temp_dir().join("audb_test_manifold_vectors.manifold");
let wal_path = temp_dir.with_extension("wal");
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
let db = Database::open(&temp_dir).await.unwrap();
let collection = db.collection("documents").unwrap();
let id = Uuid::new_v4();
let embedding = vec![0.1f32, 0.2, 0.3, 0.4];
let vectors = collection.vectors::<4>("embeddings").unwrap();
vectors.insert(&id, &embedding).unwrap();
let retrieved = vectors.get(&id).unwrap();
assert!(retrieved.is_some());
let retrieved_vec = retrieved.unwrap();
assert_eq!(retrieved_vec.len(), 4);
assert_eq!(retrieved_vec, embedding);
let wrong_dim = vec![0.1f32; 8];
let result = vectors.insert(&id, &wrong_dim);
assert!(result.is_err());
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
}
#[tokio::test]
async fn test_vector_similarity_search() {
let temp_dir = std::env::temp_dir().join("audb_test_manifold_similarity.manifold");
let wal_path = temp_dir.with_extension("wal");
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
let db = Database::open(&temp_dir).await.unwrap();
let collection = db.collection("documents").unwrap();
let vectors = collection.vectors::<4>("embeddings").unwrap();
let id1 = Uuid::new_v4();
let id2 = Uuid::new_v4();
let id3 = Uuid::new_v4();
let vec1 = vec![1.0f32, 0.0, 0.0, 0.0];
let vec2 = vec![0.7f32, 0.7, 0.0, 0.0];
let vec3 = vec![0.0f32, 0.0, 1.0, 0.0];
vectors.insert(&id1, &vec1).unwrap();
vectors.insert(&id2, &vec2).unwrap();
vectors.insert(&id3, &vec3).unwrap();
let query = vec![0.9f32, 0.1, 0.0, 0.0];
let results = vectors.search_similar(&query, 2).unwrap();
assert_eq!(results.len(), 2);
assert_eq!(results[0], id1);
let _ = std::fs::remove_file(&temp_dir);
let _ = std::fs::remove_file(&wal_path);
}
#[test]
fn test_cosine_similarity() {
let a = vec![1.0f32, 0.0, 0.0];
let b = vec![1.0f32, 0.0, 0.0];
let sim = cosine_similarity(&a, &b);
assert!((sim - 1.0).abs() < 0.001);
let a = vec![1.0f32, 0.0, 0.0];
let b = vec![0.0f32, 1.0, 0.0];
let sim = cosine_similarity(&a, &b);
assert!(sim.abs() < 0.001);
let a = vec![1.0f32, 0.0, 0.0];
let b = vec![-1.0f32, 0.0, 0.0];
let sim = cosine_similarity(&a, &b);
assert!((sim + 1.0).abs() < 0.001);
let a = vec![1.0f32, 0.5, 0.0];
let b = vec![0.9f32, 0.4, 0.0];
let sim = cosine_similarity(&a, &b);
assert!(sim > 0.99); }
}