ubiquity-database 0.1.0

Database abstraction layer for Ubiquity supporting SQLite and Astra DB
Documentation
//! Core database traits

use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use ubiquity_core::{ConsciousnessState, ConsciousnessRipple, Task, TaskResult};

/// Core database trait for all backends
#[async_trait]
pub trait Database: Send + Sync {
    /// Check database health
    async fn health_check(&self) -> Result<(), crate::DatabaseError>;
    
    /// Initialize database schema
    async fn initialize(&self) -> Result<(), crate::DatabaseError>;
    
    /// Store consciousness state
    async fn store_consciousness_state(&self, state: &ConsciousnessState) -> Result<(), crate::DatabaseError>;
    
    /// Get consciousness history for an agent
    async fn get_consciousness_history(
        &self,
        agent_id: &str,
        limit: usize,
    ) -> Result<Vec<ConsciousnessState>, crate::DatabaseError>;
    
    /// Store consciousness ripple
    async fn store_ripple(&self, ripple: &ConsciousnessRipple) -> Result<(), crate::DatabaseError>;
    
    /// Get ripples by time range
    async fn get_ripples(
        &self,
        start: DateTime<Utc>,
        end: DateTime<Utc>,
    ) -> Result<Vec<ConsciousnessRipple>, crate::DatabaseError>;
    
    /// Store task
    async fn store_task(&self, task: &Task) -> Result<(), crate::DatabaseError>;
    
    /// Update task result
    async fn store_task_result(&self, result: &TaskResult) -> Result<(), crate::DatabaseError>;
    
    /// Get pending tasks
    async fn get_pending_tasks(&self) -> Result<Vec<Task>, crate::DatabaseError>;
    
    /// Vector search for consciousness embeddings
    async fn vector_search(
        &self,
        embedding: &[f32],
        limit: usize,
    ) -> Result<Vec<VectorSearchResult>, crate::DatabaseError>;
    
    /// Hybrid search combining vector and structured queries
    async fn hybrid_search(
        &self,
        query: HybridSearchQuery,
    ) -> Result<Vec<HybridSearchResult>, crate::DatabaseError>;
}

/// Vector search result
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct VectorSearchResult {
    pub id: String,
    pub score: f32,
    pub metadata: Value,
    pub embedding: Vec<f32>,
}

/// Hybrid search query
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HybridSearchQuery {
    pub text: String,
    pub vector: Option<Vec<f32>>,
    pub filters: Option<Value>,
    pub limit: usize,
    pub rerank: bool,
}

/// Hybrid search result
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HybridSearchResult {
    pub id: String,
    pub score: f32,
    pub rerank_score: Option<f32>,
    pub metadata: Value,
    pub highlights: Vec<String>,
}

/// Memory pool for parallel processing
#[async_trait]
pub trait MemoryPool: Send + Sync {
    /// Store data in the pool
    async fn store(&self, key: &str, value: &[u8]) -> Result<(), crate::DatabaseError>;
    
    /// Retrieve data from the pool
    async fn retrieve(&self, key: &str) -> Result<Option<Vec<u8>>, crate::DatabaseError>;
    
    /// Delete data from the pool
    async fn delete(&self, key: &str) -> Result<(), crate::DatabaseError>;
    
    /// List keys with prefix
    async fn list_keys(&self, prefix: &str) -> Result<Vec<String>, crate::DatabaseError>;
    
    /// Get pool statistics
    async fn stats(&self) -> Result<PoolStats, crate::DatabaseError>;
}

/// Pool statistics
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PoolStats {
    pub size: usize,
    pub items: usize,
    pub hits: u64,
    pub misses: u64,
    pub evictions: u64,
}

/// Manager for 7-wide memory pools
#[async_trait]
pub trait MemoryPoolManager: Send + Sync {
    /// Get a specific pool by index (0-6)
    async fn get_pool(&self, index: usize) -> Result<Box<dyn MemoryPool>, crate::DatabaseError>;
    
    /// Get pool for a specific key (using consistent hashing)
    async fn get_pool_for_key(&self, key: &str) -> Result<Box<dyn MemoryPool>, crate::DatabaseError>;
    
    /// Get all pool statistics
    async fn all_stats(&self) -> Result<Vec<PoolStats>, crate::DatabaseError>;
    
    /// Rebalance pools
    async fn rebalance(&self) -> Result<(), crate::DatabaseError>;
}

/// Pub/sub messaging trait
#[async_trait]
pub trait PubSubMessaging: Send + Sync {
    /// Subscribe to a topic
    async fn subscribe(&self, topic: &str) -> Result<Box<dyn MessageSubscription>, crate::DatabaseError>;
    
    /// Publish a message
    async fn publish(&self, topic: &str, message: &[u8]) -> Result<(), crate::DatabaseError>;
    
    /// Create a topic
    async fn create_topic(&self, topic: &str) -> Result<(), crate::DatabaseError>;
    
    /// Delete a topic
    async fn delete_topic(&self, topic: &str) -> Result<(), crate::DatabaseError>;
}

/// Message subscription
#[async_trait]
pub trait MessageSubscription: Send + Sync {
    /// Receive next message
    async fn receive(&mut self) -> Result<Option<Message>, crate::DatabaseError>;
    
    /// Acknowledge message
    async fn ack(&mut self, message_id: &str) -> Result<(), crate::DatabaseError>;
    
    /// Close subscription
    async fn close(self: Box<Self>) -> Result<(), crate::DatabaseError>;
}

/// Pub/sub message
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Message {
    pub id: String,
    pub topic: String,
    pub payload: Vec<u8>,
    pub timestamp: DateTime<Utc>,
    pub properties: Option<Value>,
}