use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use ubiquity_core::{ConsciousnessState, ConsciousnessRipple, Task, TaskResult};
#[async_trait]
pub trait Database: Send + Sync {
async fn health_check(&self) -> Result<(), crate::DatabaseError>;
async fn initialize(&self) -> Result<(), crate::DatabaseError>;
async fn store_consciousness_state(&self, state: &ConsciousnessState) -> Result<(), crate::DatabaseError>;
async fn get_consciousness_history(
&self,
agent_id: &str,
limit: usize,
) -> Result<Vec<ConsciousnessState>, crate::DatabaseError>;
async fn store_ripple(&self, ripple: &ConsciousnessRipple) -> Result<(), crate::DatabaseError>;
async fn get_ripples(
&self,
start: DateTime<Utc>,
end: DateTime<Utc>,
) -> Result<Vec<ConsciousnessRipple>, crate::DatabaseError>;
async fn store_task(&self, task: &Task) -> Result<(), crate::DatabaseError>;
async fn store_task_result(&self, result: &TaskResult) -> Result<(), crate::DatabaseError>;
async fn get_pending_tasks(&self) -> Result<Vec<Task>, crate::DatabaseError>;
async fn vector_search(
&self,
embedding: &[f32],
limit: usize,
) -> Result<Vec<VectorSearchResult>, crate::DatabaseError>;
async fn hybrid_search(
&self,
query: HybridSearchQuery,
) -> Result<Vec<HybridSearchResult>, crate::DatabaseError>;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct VectorSearchResult {
pub id: String,
pub score: f32,
pub metadata: Value,
pub embedding: Vec<f32>,
}
#[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,
}
#[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>,
}
#[async_trait]
pub trait MemoryPool: Send + Sync {
async fn store(&self, key: &str, value: &[u8]) -> Result<(), crate::DatabaseError>;
async fn retrieve(&self, key: &str) -> Result<Option<Vec<u8>>, crate::DatabaseError>;
async fn delete(&self, key: &str) -> Result<(), crate::DatabaseError>;
async fn list_keys(&self, prefix: &str) -> Result<Vec<String>, crate::DatabaseError>;
async fn stats(&self) -> Result<PoolStats, crate::DatabaseError>;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PoolStats {
pub size: usize,
pub items: usize,
pub hits: u64,
pub misses: u64,
pub evictions: u64,
}
#[async_trait]
pub trait MemoryPoolManager: Send + Sync {
async fn get_pool(&self, index: usize) -> Result<Box<dyn MemoryPool>, crate::DatabaseError>;
async fn get_pool_for_key(&self, key: &str) -> Result<Box<dyn MemoryPool>, crate::DatabaseError>;
async fn all_stats(&self) -> Result<Vec<PoolStats>, crate::DatabaseError>;
async fn rebalance(&self) -> Result<(), crate::DatabaseError>;
}
#[async_trait]
pub trait PubSubMessaging: Send + Sync {
async fn subscribe(&self, topic: &str) -> Result<Box<dyn MessageSubscription>, crate::DatabaseError>;
async fn publish(&self, topic: &str, message: &[u8]) -> Result<(), crate::DatabaseError>;
async fn create_topic(&self, topic: &str) -> Result<(), crate::DatabaseError>;
async fn delete_topic(&self, topic: &str) -> Result<(), crate::DatabaseError>;
}
#[async_trait]
pub trait MessageSubscription: Send + Sync {
async fn receive(&mut self) -> Result<Option<Message>, crate::DatabaseError>;
async fn ack(&mut self, message_id: &str) -> Result<(), crate::DatabaseError>;
async fn close(self: Box<Self>) -> Result<(), crate::DatabaseError>;
}
#[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>,
}