Skip to main content

ubiquity_database/
traits.rs

1//! Core database traits
2
3use async_trait::async_trait;
4use chrono::{DateTime, Utc};
5use serde::{Deserialize, Serialize};
6use serde_json::Value;
7use ubiquity_core::{ConsciousnessState, ConsciousnessRipple, Task, TaskResult};
8
9/// Core database trait for all backends
10#[async_trait]
11pub trait Database: Send + Sync {
12    /// Check database health
13    async fn health_check(&self) -> Result<(), crate::DatabaseError>;
14    
15    /// Initialize database schema
16    async fn initialize(&self) -> Result<(), crate::DatabaseError>;
17    
18    /// Store consciousness state
19    async fn store_consciousness_state(&self, state: &ConsciousnessState) -> Result<(), crate::DatabaseError>;
20    
21    /// Get consciousness history for an agent
22    async fn get_consciousness_history(
23        &self,
24        agent_id: &str,
25        limit: usize,
26    ) -> Result<Vec<ConsciousnessState>, crate::DatabaseError>;
27    
28    /// Store consciousness ripple
29    async fn store_ripple(&self, ripple: &ConsciousnessRipple) -> Result<(), crate::DatabaseError>;
30    
31    /// Get ripples by time range
32    async fn get_ripples(
33        &self,
34        start: DateTime<Utc>,
35        end: DateTime<Utc>,
36    ) -> Result<Vec<ConsciousnessRipple>, crate::DatabaseError>;
37    
38    /// Store task
39    async fn store_task(&self, task: &Task) -> Result<(), crate::DatabaseError>;
40    
41    /// Update task result
42    async fn store_task_result(&self, result: &TaskResult) -> Result<(), crate::DatabaseError>;
43    
44    /// Get pending tasks
45    async fn get_pending_tasks(&self) -> Result<Vec<Task>, crate::DatabaseError>;
46    
47    /// Vector search for consciousness embeddings
48    async fn vector_search(
49        &self,
50        embedding: &[f32],
51        limit: usize,
52    ) -> Result<Vec<VectorSearchResult>, crate::DatabaseError>;
53    
54    /// Hybrid search combining vector and structured queries
55    async fn hybrid_search(
56        &self,
57        query: HybridSearchQuery,
58    ) -> Result<Vec<HybridSearchResult>, crate::DatabaseError>;
59}
60
61/// Vector search result
62#[derive(Debug, Clone, Serialize, Deserialize)]
63pub struct VectorSearchResult {
64    pub id: String,
65    pub score: f32,
66    pub metadata: Value,
67    pub embedding: Vec<f32>,
68}
69
70/// Hybrid search query
71#[derive(Debug, Clone, Serialize, Deserialize)]
72pub struct HybridSearchQuery {
73    pub text: String,
74    pub vector: Option<Vec<f32>>,
75    pub filters: Option<Value>,
76    pub limit: usize,
77    pub rerank: bool,
78}
79
80/// Hybrid search result
81#[derive(Debug, Clone, Serialize, Deserialize)]
82pub struct HybridSearchResult {
83    pub id: String,
84    pub score: f32,
85    pub rerank_score: Option<f32>,
86    pub metadata: Value,
87    pub highlights: Vec<String>,
88}
89
90/// Memory pool for parallel processing
91#[async_trait]
92pub trait MemoryPool: Send + Sync {
93    /// Store data in the pool
94    async fn store(&self, key: &str, value: &[u8]) -> Result<(), crate::DatabaseError>;
95    
96    /// Retrieve data from the pool
97    async fn retrieve(&self, key: &str) -> Result<Option<Vec<u8>>, crate::DatabaseError>;
98    
99    /// Delete data from the pool
100    async fn delete(&self, key: &str) -> Result<(), crate::DatabaseError>;
101    
102    /// List keys with prefix
103    async fn list_keys(&self, prefix: &str) -> Result<Vec<String>, crate::DatabaseError>;
104    
105    /// Get pool statistics
106    async fn stats(&self) -> Result<PoolStats, crate::DatabaseError>;
107}
108
109/// Pool statistics
110#[derive(Debug, Clone, Serialize, Deserialize)]
111pub struct PoolStats {
112    pub size: usize,
113    pub items: usize,
114    pub hits: u64,
115    pub misses: u64,
116    pub evictions: u64,
117}
118
119/// Manager for 7-wide memory pools
120#[async_trait]
121pub trait MemoryPoolManager: Send + Sync {
122    /// Get a specific pool by index (0-6)
123    async fn get_pool(&self, index: usize) -> Result<Box<dyn MemoryPool>, crate::DatabaseError>;
124    
125    /// Get pool for a specific key (using consistent hashing)
126    async fn get_pool_for_key(&self, key: &str) -> Result<Box<dyn MemoryPool>, crate::DatabaseError>;
127    
128    /// Get all pool statistics
129    async fn all_stats(&self) -> Result<Vec<PoolStats>, crate::DatabaseError>;
130    
131    /// Rebalance pools
132    async fn rebalance(&self) -> Result<(), crate::DatabaseError>;
133}
134
135/// Pub/sub messaging trait
136#[async_trait]
137pub trait PubSubMessaging: Send + Sync {
138    /// Subscribe to a topic
139    async fn subscribe(&self, topic: &str) -> Result<Box<dyn MessageSubscription>, crate::DatabaseError>;
140    
141    /// Publish a message
142    async fn publish(&self, topic: &str, message: &[u8]) -> Result<(), crate::DatabaseError>;
143    
144    /// Create a topic
145    async fn create_topic(&self, topic: &str) -> Result<(), crate::DatabaseError>;
146    
147    /// Delete a topic
148    async fn delete_topic(&self, topic: &str) -> Result<(), crate::DatabaseError>;
149}
150
151/// Message subscription
152#[async_trait]
153pub trait MessageSubscription: Send + Sync {
154    /// Receive next message
155    async fn receive(&mut self) -> Result<Option<Message>, crate::DatabaseError>;
156    
157    /// Acknowledge message
158    async fn ack(&mut self, message_id: &str) -> Result<(), crate::DatabaseError>;
159    
160    /// Close subscription
161    async fn close(self: Box<Self>) -> Result<(), crate::DatabaseError>;
162}
163
164/// Pub/sub message
165#[derive(Debug, Clone, Serialize, Deserialize)]
166pub struct Message {
167    pub id: String,
168    pub topic: String,
169    pub payload: Vec<u8>,
170    pub timestamp: DateTime<Utc>,
171    pub properties: Option<Value>,
172}